@cryptotaxi247 / netdata-1 / commits / 95119afff

Make dbengine the default memory mode (#6977)

* Basic functionality for dbengine stress test. * Fix coverity defects * Refactored dbengine stress test to be configurable * Added benchmark results and evaluation in dbengine documentation * Make dbengine the default memory mode

Markos Fountoulakis committed Oct 3, 2019 at 17:04 UTC 95119afff48735607643bfe3824ed3727b6edbb0
8 files changed +376 -88
daemon/main.c
+29 -3
@@ -306,7 +306,13 @@ int help(int exitcode) {
306 " -W stacksize=N Set the stacksize (in bytes).\n\n"
307 " -W debug_flags=N Set runtime tracing to debug.log.\n\n"
308 " -W unittest Run internal unittests and exit.\n\n"
309 +#ifdef ENABLE_DBENGINE
310 " -W createdataset=N Create a DB engine dataset of N seconds and exit.\n\n"
311 + " -W stresstest=A,B,C,D,E Run a DB engine stress test for A seconds,\n"
312 + " with B writers and C readers, with a ramp up\n"
313 + " time of D seconds for writers, a page cache\n"
314 + " size of E MiB, and exit.\n\n"
315 +#endif
316 " -W set section option value\n"
317 " set netdata.conf option from the command line.\n\n"
318 " -W simple-pattern pattern string\n"
@@ -887,6 +893,7 @@ int main(int argc, char **argv) {
893 char* stacksize_string = "stacksize=";
894 char* debug_flags_string = "debug_flags=";
895 char* createdataset_string = "createdataset=";
896 + char* stresstest_string = "stresstest=";
897
898 if(strcmp(optarg, "unittest") == 0) {
899 if(unit_test_buffer()) return 1;
@@ -905,14 +912,33 @@ int main(int argc, char **argv) {
912 fprintf(stderr, "\n\nALL TESTS PASSED\n\n");
913 return 0;
914 }
915 +#ifdef ENABLE_DBENGINE
916 else if(strncmp(optarg, createdataset_string, strlen(createdataset_string)) == 0) {
917 optarg += strlen(createdataset_string);
910 -#ifdef ENABLE_DBENGINE
911 - unsigned history_seconds = (unsigned )strtoull(optarg, NULL, 0);
918 + unsigned history_seconds = strtoul(optarg, NULL, 0);
919 generate_dbengine_dataset(history_seconds);
913 -#endif
920 return 0;
921 }
922 + else if(strncmp(optarg, stresstest_string, strlen(stresstest_string)) == 0) {
923 + char *endptr;
924 + unsigned test_duration_sec = 0, dset_charts = 0, query_threads = 0, ramp_up_seconds = 0,
925 + page_cache_mb = 0;
926 +
927 + optarg += strlen(stresstest_string);
928 + test_duration_sec = (unsigned)strtoul(optarg, &endptr, 0);
929 + if (',' == *endptr)
930 + dset_charts = (unsigned)strtoul(endptr + 1, &endptr, 0);
931 + if (',' == *endptr)
932 + query_threads = (unsigned)strtoul(endptr + 1, &endptr, 0);
933 + if (',' == *endptr)
934 + ramp_up_seconds = (unsigned)strtoul(endptr + 1, &endptr, 0);
935 + if (',' == *endptr)
936 + page_cache_mb = (unsigned)strtoul(endptr + 1, &endptr, 0);
937 + dbengine_stress_test(test_duration_sec, dset_charts, query_threads, ramp_up_seconds,
938 + page_cache_mb);
939 + return 0;
940 + }
941 +#endif
942 else if(strcmp(optarg, "simple-pattern") == 0) {
943 if(optind + 2 > argc) {
944 fprintf(stderr, "%s", "\nUSAGE: -W simple-pattern 'pattern' 'string'\n\n"
daemon/unit_test.c
+279 -33
@@ -1688,7 +1688,8 @@ static time_t test_dbengine_create_metrics(RRDSET *st[CHARTS], RRDDIM *rd[CHARTS
1688 st[i]->usec_since_last_update = USEC_PER_SEC * update_every;
1689
1690 for (j = 0; j < DIMS; ++j) {
1691 - next = i * DIMS * REGION_POINTS[current_region] + j * REGION_POINTS[current_region] + c;
1691 + next = ((collected_number)i * DIMS) * REGION_POINTS[current_region] +
1692 + j * REGION_POINTS[current_region] + c;
1693 rrddim_set_by_pointer_fake_time(rd[i][j], next, time_now);
1694 }
1695 rrdset_done(st[i]);
@@ -1719,13 +1720,14 @@ static int test_dbengine_check_metrics(RRDSET *st[CHARTS], RRDDIM *rd[CHARTS][DI
1720 for (j = 0; j < DIMS; ++j) {
1721 rd[i][j]->state->query_ops.init(rd[i][j], &handle, time_now, time_now + QUERY_BATCH * update_every);
1722 for (k = 0; k < QUERY_BATCH; ++k) {
1722 - last = i * DIMS * REGION_POINTS[current_region] + j * REGION_POINTS[current_region] + c + k;
1723 + last = ((collected_number)i * DIMS) * REGION_POINTS[current_region] +
1724 + j * REGION_POINTS[current_region] + c + k;
1725 expected = unpack_storage_number(pack_storage_number((calculated_number)last, SN_EXISTS));
1726
1727 n = rd[i][j]->state->query_ops.next_metric(&handle, &time_retrieved);
1728 value = unpack_storage_number(n);
1729
1728 - same = (calculated_number_round(value * 10000000.0) == calculated_number_round(expected * 10000000.0)) ? 1 : 0;
1730 + same = (calculated_number_round(value) == calculated_number_round(expected)) ? 1 : 0;
1731 if(!same) {
1732 fprintf(stderr, " DB-engine unittest %s/%s: at %lu secs, expecting value "
1733 CALCULATED_NUMBER_FORMAT ", found " CALCULATED_NUMBER_FORMAT ", ### E R R O R ###\n",
@@ -1780,7 +1782,7 @@ static int test_dbengine_check_rrdr(RRDSET *st[CHARTS], RRDDIM *rd[CHARTS][DIMS]
1782 last = i * DIMS * REGION_POINTS[current_region] + j * REGION_POINTS[current_region] + c;
1783 expected = unpack_storage_number(pack_storage_number((calculated_number)last, SN_EXISTS));
1784
1783 - same = (calculated_number_round(value * 10000000.0) == calculated_number_round(expected * 10000000.0)) ? 1 : 0;
1785 + same = (calculated_number_round(value) == calculated_number_round(expected)) ? 1 : 0;
1786 if(!same) {
1787 fprintf(stderr, " DB-engine unittest %s/%s: at %lu secs, expecting value "
1788 CALCULATED_NUMBER_FORMAT ", RRDR found " CALCULATED_NUMBER_FORMAT ", ### E R R O R ###\n",
@@ -1902,7 +1904,7 @@ int test_dbengine(void)
1904 collected_number last = i * DIMS * REGION_POINTS[current_region] + j * REGION_POINTS[current_region] + c - point_offset;
1905 calculated_number expected = unpack_storage_number(pack_storage_number((calculated_number)last, SN_EXISTS));
1906
1905 - uint8_t same = (calculated_number_round(value * 10000000.0) == calculated_number_round(expected * 10000000.0)) ? 1 : 0;
1907 + uint8_t same = (calculated_number_round(value) == calculated_number_round(expected)) ? 1 : 0;
1908 if(!same) {
1909 fprintf(stderr, " DB-engine unittest %s/%s: at %lu secs, expecting value "
1910 CALCULATED_NUMBER_FORMAT ", RRDR found " CALCULATED_NUMBER_FORMAT ", ### E R R O R ###\n",
@@ -1932,20 +1934,27 @@ struct dbengine_chart_thread {
1934 uv_thread_t thread;
1935 RRDHOST *host;
1936 char *chartname; /* Will be prefixed by type, e.g. "example_local1.", "example_local2." etc */
1935 - int dset_charts; /* number of charts */
1936 - int dset_dims; /* dimensions per chart */
1937 - int chart_i; /* current chart offset */
1937 + unsigned dset_charts; /* number of charts */
1938 + unsigned dset_dims; /* dimensions per chart */
1939 + unsigned chart_i; /* current chart offset */
1940 time_t time_present; /* current virtual time of the benchmark */
1941 + volatile time_t time_max; /* latest timestamp of stored values */
1942 unsigned history_seconds; /* how far back in the past to go */
1943 +
1944 + volatile long done; /* initialize to 0, set to 1 to stop thread */
1945 + struct completion charts_initialized;
1946 + unsigned long errors, stored_metrics_nr; /* statistics */
1947 +
1948 + RRDSET *st;
1949 + RRDDIM *rd[]; /* dset_dims elements */
1950 };
1951
1942 -collected_number generate_dbengine_chart_value(struct dbengine_chart_thread *thread_info, int dim_i,
1943 - time_t time_current)
1952 +collected_number generate_dbengine_chart_value(int chart_i, int dim_i, time_t time_current)
1953 {
1954 collected_number value;
1955
1947 - value = ((collected_number)time_current) * thread_info->chart_i;
1948 - value += ((collected_number)time_current) * dim_i;
1956 + value = ((collected_number)time_current) * (chart_i + 1);
1957 + value += ((collected_number)time_current) * (dim_i + 1);
1958 value %= 1024LLU;
1959
1960 return value;
@@ -1956,44 +1965,47 @@ static void generate_dbengine_chart(void *arg)
1965 struct dbengine_chart_thread *thread_info = (struct dbengine_chart_thread *)arg;
1966 RRDHOST *host = thread_info->host;
1967 char *chartname = thread_info->chartname;
1959 - const int DSET_DIMS = thread_info->dset_dims;
1968 + const unsigned DSET_DIMS = thread_info->dset_dims;
1969 unsigned history_seconds = thread_info->history_seconds;
1970 time_t time_present = thread_info->time_present;
1971
1963 - int j, update_every = 1;
1972 + unsigned j, update_every = 1;
1973 RRDSET *st;
1974 RRDDIM *rd[DSET_DIMS];
1975 char name[RRD_ID_LENGTH_MAX + 1];
1976 time_t time_current;
1977
1978 // create the chart
1970 - snprintfz(name, RRD_ID_LENGTH_MAX, "example_local%d", thread_info->chart_i + 1);
1971 - st = rrdset_create(host, name, chartname, chartname, "example", NULL, chartname, chartname, chartname, NULL, 1,
1972 - update_every, RRDSET_TYPE_LINE);
1979 + snprintfz(name, RRD_ID_LENGTH_MAX, "example_local%u", thread_info->chart_i + 1);
1980 + thread_info->st = st = rrdset_create(host, name, chartname, chartname, "example", NULL, chartname, chartname,
1981 + chartname, NULL, 1, update_every, RRDSET_TYPE_LINE);
1982 for (j = 0 ; j < DSET_DIMS ; ++j) {
1974 - snprintfz(name, RRD_ID_LENGTH_MAX, "%s%d", chartname, j);
1983 + snprintfz(name, RRD_ID_LENGTH_MAX, "%s%u", chartname, j + 1);
1984
1976 - rd[j] = rrddim_add(st, name, NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
1985 + thread_info->rd[j] = rd[j] = rrddim_add(st, name, NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
1986 }
1987 + complete(&thread_info->charts_initialized);
1988
1989 // feed it with the test data
1990 time_current = time_present - history_seconds;
1991 for (j = 0 ; j < DSET_DIMS ; ++j) {
1992 rd[j]->last_collected_time.tv_sec =
1983 - st->last_collected_time.tv_sec = st->last_updated.tv_sec = time_current;
1993 + st->last_collected_time.tv_sec = st->last_updated.tv_sec = time_current - update_every;
1994 rd[j]->last_collected_time.tv_usec =
1995 st->last_collected_time.tv_usec = st->last_updated.tv_usec = 0;
1996 }
1987 - for( ; time_current < time_present ; ++time_current) {
1988 - st->usec_since_last_update = USEC_PER_SEC;
1997 + for( ; !thread_info->done && time_current < time_present ; time_current += update_every) {
1998 + st->usec_since_last_update = USEC_PER_SEC * update_every;
1999
2000 for (j = 0; j < DSET_DIMS; ++j) {
2001 collected_number value;
2002
1993 - value = generate_dbengine_chart_value(thread_info, j, time_current);
2003 + value = generate_dbengine_chart_value(thread_info->chart_i, j, time_current);
2004 rrddim_set_by_pointer_fake_time(rd[j], value, time_current);
2005 + ++thread_info->stored_metrics_nr;
2006 }
2007 rrdset_done(st);
2008 + thread_info->time_max = time_current;
2009 }
2010 }
2011
@@ -2003,7 +2015,7 @@ void generate_dbengine_dataset(unsigned history_seconds)
2015 const int DSET_DIMS = 128;
2016 const uint64_t EXPECTED_COMPRESSION_RATIO = 20;
2017 RRDHOST *host = NULL;
2006 - struct dbengine_chart_thread thread_info[DSET_CHARTS];
2018 + struct dbengine_chart_thread **thread_info;
2019 int i;
2020 time_t time_present;
2021
@@ -2021,25 +2033,259 @@ void generate_dbengine_dataset(unsigned history_seconds)
2033 if (NULL == host)
2034 return;
2035
2036 + thread_info = mallocz(sizeof(*thread_info) * DSET_CHARTS);
2037 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
2038 + thread_info[i] = mallocz(sizeof(*thread_info[i]) + sizeof(RRDDIM *) * DSET_DIMS);
2039 + }
2040 fprintf(stderr, "\nRunning DB-engine workload generator\n");
2041
2042 time_present = now_realtime_sec();
2043 for (i = 0 ; i < DSET_CHARTS ; ++i) {
2028 - thread_info[i].host = host;
2029 - thread_info[i].chartname = "random";
2030 - thread_info[i].dset_charts = DSET_CHARTS;
2031 - thread_info[i].chart_i = i;
2032 - thread_info[i].dset_dims = DSET_DIMS;
2033 - thread_info[i].history_seconds = history_seconds;
2034 - thread_info[i].time_present = time_present;
2035 - assert(0 == uv_thread_create(&thread_info[i].thread, generate_dbengine_chart, &thread_info[i]));
2044 + thread_info[i]->host = host;
2045 + thread_info[i]->chartname = "random";
2046 + thread_info[i]->dset_charts = DSET_CHARTS;
2047 + thread_info[i]->chart_i = i;
2048 + thread_info[i]->dset_dims = DSET_DIMS;
2049 + thread_info[i]->history_seconds = history_seconds;
2050 + thread_info[i]->time_present = time_present;
2051 + thread_info[i]->time_max = 0;
2052 + thread_info[i]->done = 0;
2053 + init_completion(&thread_info[i]->charts_initialized);
2054 + assert(0 == uv_thread_create(&thread_info[i]->thread, generate_dbengine_chart, thread_info[i]));
2055 + wait_for_completion(&thread_info[i]->charts_initialized);
2056 + destroy_completion(&thread_info[i]->charts_initialized);
2057 }
2058 for (i = 0 ; i < DSET_CHARTS ; ++i) {
2038 - assert(0 == uv_thread_join(&thread_info[i].thread));
2059 + assert(0 == uv_thread_join(&thread_info[i]->thread));
2060 }
2061
2062 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
2063 + freez(thread_info[i]);
2064 + }
2065 + freez(thread_info);
2066 rrd_wrlock();
2067 rrdhost_free(host);
2068 rrd_unlock();
2069 }
2070 +
2071 +struct dbengine_query_thread {
2072 + uv_thread_t thread;
2073 + RRDHOST *host;
2074 + char *chartname; /* Will be prefixed by type, e.g. "example_local1.", "example_local2." etc */
2075 + unsigned dset_charts; /* number of charts */
2076 + unsigned dset_dims; /* dimensions per chart */
2077 + time_t time_present; /* current virtual time of the benchmark */
2078 + unsigned history_seconds; /* how far back in the past to go */
2079 + volatile long done; /* initialize to 0, set to 1 to stop thread */
2080 + unsigned long errors, queries_nr, queried_metrics_nr; /* statistics */
2081 +
2082 + struct dbengine_chart_thread *chart_threads[]; /* dset_charts elements */
2083 +};
2084 +
2085 +static void query_dbengine_chart(void *arg)
2086 +{
2087 + struct dbengine_query_thread *thread_info = (struct dbengine_query_thread *)arg;
2088 + const int DSET_CHARTS = thread_info->dset_charts;
2089 + const int DSET_DIMS = thread_info->dset_dims;
2090 + time_t time_after, time_before, time_min, time_max, duration;
2091 + int i, j, update_every = 1;
2092 + RRDSET *st;
2093 + RRDDIM *rd;
2094 + uint8_t same;
2095 + time_t time_now, time_retrieved;
2096 + collected_number generatedv;
2097 + calculated_number value, expected;
2098 + storage_number n;
2099 + struct rrddim_query_handle handle;
2100 +
2101 + do {
2102 + // pick a chart and dimension
2103 + i = random() % DSET_CHARTS;
2104 + st = thread_info->chart_threads[i]->st;
2105 + j = random() % DSET_DIMS;
2106 + rd = thread_info->chart_threads[i]->rd[j];
2107 +
2108 + time_min = thread_info->time_present - thread_info->history_seconds + 1;
2109 + time_max = thread_info->chart_threads[i]->time_max;
2110 + if (!time_max) {
2111 + time_before = time_after = time_min;
2112 + } else {
2113 + time_after = time_min + random() % (MAX(time_max - time_min, 1));
2114 + duration = random() % 3600;
2115 + time_before = MIN(time_after + duration, time_max); /* up to 1 hour queries */
2116 + }
2117 +
2118 + rd->state->query_ops.init(rd, &handle, time_after, time_before);
2119 + ++thread_info->queries_nr;
2120 + for (time_now = time_after ; time_now <= time_before ; time_now += update_every) {
2121 + generatedv = generate_dbengine_chart_value(i, j, time_now);
2122 + expected = unpack_storage_number(pack_storage_number((calculated_number) generatedv, SN_EXISTS));
2123 +
2124 + if (unlikely(rd->state->query_ops.is_finished(&handle))) {
2125 + fprintf(stderr, " DB-engine stresstest %s/%s: at %lu secs, expecting value "
2126 + CALCULATED_NUMBER_FORMAT ", found data gap, ### E R R O R ###\n",
2127 + st->name, rd->name, (unsigned long) time_now, expected);
2128 + ++thread_info->errors;
2129 + break;
2130 + }
2131 + n = rd->state->query_ops.next_metric(&handle, &time_retrieved);
2132 + if (SN_EMPTY_SLOT == n) {
2133 + fprintf(stderr, " DB-engine stresstest %s/%s: at %lu secs, expecting value "
2134 + CALCULATED_NUMBER_FORMAT ", found data gap, ### E R R O R ###\n",
2135 + st->name, rd->name, (unsigned long) time_now, expected);
2136 + ++thread_info->errors;
2137 + break;
2138 + }
2139 + ++thread_info->queried_metrics_nr;
2140 + value = unpack_storage_number(n);
2141 +
2142 + same = (calculated_number_round(value) == calculated_number_round(expected)) ? 1 : 0;
2143 + if (!same) {
2144 + fprintf(stderr, " DB-engine stresstest %s/%s: at %lu secs, expecting value "
2145 + CALCULATED_NUMBER_FORMAT ", found " CALCULATED_NUMBER_FORMAT ", ### E R R O R ###\n",
2146 + st->name, rd->name, (unsigned long) time_now, expected, value);
2147 + ++thread_info->errors;
2148 + }
2149 + if (time_retrieved != time_now) {
2150 + fprintf(stderr, " DB-engine stresstest %s/%s: at %lu secs, found timestamp %lu ### E R R O R ###\n",
2151 + st->name, rd->name, (unsigned long) time_now, (unsigned long) time_retrieved);
2152 + ++thread_info->errors;
2153 + }
2154 + }
2155 + rd->state->query_ops.finalize(&handle);
2156 + } while(!thread_info->done);
2157 +}
2158 +
2159 +void dbengine_stress_test(unsigned TEST_DURATION_SEC, unsigned DSET_CHARTS, unsigned QUERY_THREADS,
2160 + unsigned RAMP_UP_SECONDS, unsigned PAGE_CACHE_MB)
2161 +{
2162 + const unsigned DSET_DIMS = 128;
2163 + const uint64_t EXPECTED_COMPRESSION_RATIO = 20;
2164 + const unsigned HISTORY_SECONDS = 3600 * 24 * 365; /* 1 year of history */
2165 + RRDHOST *host = NULL;
2166 + struct dbengine_chart_thread **chart_threads;
2167 + struct dbengine_query_thread **query_threads;
2168 + unsigned i, j;
2169 + time_t time_start, time_end;
2170 +
2171 + if (!TEST_DURATION_SEC)
2172 + TEST_DURATION_SEC = 10;
2173 + if (!DSET_CHARTS)
2174 + DSET_CHARTS = 1;
2175 + if (!QUERY_THREADS)
2176 + QUERY_THREADS = 1;
2177 + if (PAGE_CACHE_MB < RRDENG_MIN_PAGE_CACHE_SIZE_MB)
2178 + PAGE_CACHE_MB = RRDENG_MIN_PAGE_CACHE_SIZE_MB;
2179 +
2180 + default_rrd_memory_mode = RRD_MEMORY_MODE_DBENGINE;
2181 + default_rrdeng_page_cache_mb = PAGE_CACHE_MB;
2182 + // Worst case for uncompressible data
2183 + default_rrdeng_disk_quota_mb = (((uint64_t)DSET_DIMS * DSET_CHARTS) * sizeof(storage_number) * HISTORY_SECONDS) /
2184 + (1024 * 1024);
2185 + default_rrdeng_disk_quota_mb -= default_rrdeng_disk_quota_mb * EXPECTED_COMPRESSION_RATIO / 100;
2186 +
2187 + error_log_limit_unlimited();
2188 + debug(D_RRDHOST, "Initializing localhost with hostname 'dbengine-stress-test'");
2189 +
2190 + host = dbengine_rrdhost_find_or_create("dbengine-stress-test");
2191 + if (NULL == host)
2192 + return;
2193 +
2194 + chart_threads = mallocz(sizeof(*chart_threads) * DSET_CHARTS);
2195 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
2196 + chart_threads[i] = mallocz(sizeof(*chart_threads[i]) + sizeof(RRDDIM *) * DSET_DIMS);
2197 + }
2198 + query_threads = mallocz(sizeof(*query_threads) * QUERY_THREADS);
2199 + for (i = 0 ; i < QUERY_THREADS ; ++i) {
2200 + query_threads[i] = mallocz(sizeof(*query_threads[i]) + sizeof(struct dbengine_chart_thread *) * DSET_CHARTS);
2201 + }
2202 + fprintf(stderr, "\nRunning DB-engine stress test, %u seconds writers ramp-up time,\n"
2203 + "%u seconds of concurrent readers and writers, %u writer threads, %u reader threads,\n"
2204 + "%u MiB of page cache.\n",
2205 + RAMP_UP_SECONDS, TEST_DURATION_SEC, DSET_CHARTS, QUERY_THREADS, PAGE_CACHE_MB);
2206 +
2207 + time_start = now_realtime_sec();
2208 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
2209 + chart_threads[i]->host = host;
2210 + chart_threads[i]->chartname = "random";
2211 + chart_threads[i]->dset_charts = DSET_CHARTS;
2212 + chart_threads[i]->chart_i = i;
2213 + chart_threads[i]->dset_dims = DSET_DIMS;
2214 + chart_threads[i]->history_seconds = HISTORY_SECONDS;
2215 + chart_threads[i]->time_present = time_start;
2216 + chart_threads[i]->time_max = 0;
2217 + chart_threads[i]->done = 0;
2218 + chart_threads[i]->errors = chart_threads[i]->stored_metrics_nr = 0;
2219 + init_completion(&chart_threads[i]->charts_initialized);
2220 + assert(0 == uv_thread_create(&chart_threads[i]->thread, generate_dbengine_chart, chart_threads[i]));
2221 + }
2222 + /* barrier so that subsequent queries can access valid chart data */
2223 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
2224 + wait_for_completion(&chart_threads[i]->charts_initialized);
2225 + destroy_completion(&chart_threads[i]->charts_initialized);
2226 + }
2227 + sleep(RAMP_UP_SECONDS);
2228 + /* at this point data have already began being written to the database */
2229 + for (i = 0 ; i < QUERY_THREADS ; ++i) {
2230 + query_threads[i]->host = host;
2231 + query_threads[i]->chartname = "random";
2232 + query_threads[i]->dset_charts = DSET_CHARTS;
2233 + query_threads[i]->dset_dims = DSET_DIMS;
2234 + query_threads[i]->history_seconds = HISTORY_SECONDS;
2235 + query_threads[i]->time_present = time_start;
2236 + query_threads[i]->done = 0;
2237 + query_threads[i]->errors = query_threads[i]->queries_nr = query_threads[i]->queried_metrics_nr = 0;
2238 + for (j = 0 ; j < DSET_CHARTS ; ++j) {
2239 + query_threads[i]->chart_threads[j] = chart_threads[j];
2240 + }
2241 + assert(0 == uv_thread_create(&query_threads[i]->thread, query_dbengine_chart, query_threads[i]));
2242 + }
2243 + sleep(TEST_DURATION_SEC);
2244 + /* stop workload */
2245 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
2246 + chart_threads[i]->done = 1;
2247 + }
2248 + for (i = 0 ; i < QUERY_THREADS ; ++i) {
2249 + query_threads[i]->done = 1;
2250 + }
2251 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
2252 + assert(0 == uv_thread_join(&chart_threads[i]->thread));
2253 + }
2254 + for (i = 0 ; i < QUERY_THREADS ; ++i) {
2255 + assert(0 == uv_thread_join(&query_threads[i]->thread));
2256 + }
2257 + time_end = now_realtime_sec();
2258 + fprintf(stderr, "\nDB-engine stress test finished in %ld seconds.\n", time_end - time_start);
2259 + unsigned long stored_metrics_nr = 0;
2260 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
2261 + stored_metrics_nr += chart_threads[i]->stored_metrics_nr;
2262 + }
2263 + unsigned long queries_nr = 0, queried_metrics_nr = 0;
2264 + for (i = 0 ; i < QUERY_THREADS ; ++i) {
2265 + queries_nr += query_threads[i]->queries_nr;
2266 + queried_metrics_nr += query_threads[i]->queried_metrics_nr;
2267 + }
2268 + fprintf(stderr, "%u metrics were stored (dataset size of %lu MiB) in %u charts by 1 writer thread per chart.\n",
2269 + DSET_CHARTS * DSET_DIMS, stored_metrics_nr * sizeof(storage_number) / (1024 * 1024), DSET_CHARTS);
2270 + fprintf(stderr, "Metrics were being generated per 1 emulated second and time was accelerated.\n");
2271 + fprintf(stderr, "%lu metric data points were queried by %u reader threads.\n", queried_metrics_nr, QUERY_THREADS);
2272 + fprintf(stderr, "Query starting time is randomly chosen from the beginning of the time-series up to the time of\n"
2273 + "the latest data point, and ending time from 1 second up to 1 hour after the starting time.\n");
2274 + fprintf(stderr, "Performance is %lu written data points/sec and %lu read data points/sec.\n",
2275 + stored_metrics_nr / (time_end - time_start), queried_metrics_nr / (time_end - time_start));
2276 +
2277 + for (i = 0 ; i < DSET_CHARTS ; ++i) {
2278 + freez(chart_threads[i]);
2279 + }
2280 + freez(chart_threads);
2281 + for (i = 0 ; i < QUERY_THREADS ; ++i) {
2282 + freez(query_threads[i]);
2283 + }
2284 + freez(query_threads);
2285 + rrdeng_exit(host->rrdeng_ctx);
2286 + rrd_wrlock();
2287 + rrdhost_delete_charts(host);
2288 + rrd_unlock();
2289 +}
2290 +
2291 #endif
daemon/unit_test.h
+3
@@ -11,6 +11,9 @@ extern int unit_test_buffer(void);
11 #ifdef ENABLE_DBENGINE
12 extern int test_dbengine(void);
13 extern void generate_dbengine_dataset(unsigned history_seconds);
14 +extern void dbengine_stress_test(unsigned TEST_DURATION_SEC, unsigned DSET_CHARTS, unsigned QUERY_THREADS,
15 + unsigned RAMP_UP_SECONDS, unsigned PAGE_CACHE_MB);
16 +
17 #endif
18
19 #endif /* NETDATA_UNIT_TEST_H */
database/README.md
+8 -7
@@ -25,7 +25,7 @@ Currently Netdata supports 6 memory modes:
25
26 1. `ram`, data are purely in memory. Data are never saved on disk. This mode uses `mmap()` and supports [KSM](#ksm).
27
28 -2. `save`, (the default) data are only in RAM while Netdata runs and are saved to / loaded from disk on Netdata
28 +2. `save`, data are only in RAM while Netdata runs and are saved to / loaded from disk on Netdata
29 restart. It also uses `mmap()` and supports [KSM](#ksm).
30
31 3. `map`, data are in memory mapped files. This works like the swap. Keep in mind though, this will have a constant
@@ -39,11 +39,12 @@ Currently Netdata supports 6 memory modes:
39 5. `alloc`, like `ram` but it uses `calloc()` and does not support [KSM](#ksm). This mode is the fallback for all
40 others except `none`.
41
42 -6. `dbengine`, data are in database files. The [Database Engine](engine/) works like a traditional database. There is
43 - some amount of RAM dedicated to data caching and indexing and the rest of the data reside compressed on disk. The
44 - number of history entries is not fixed in this case, but depends on the configured disk space and the effective
45 - compression ratio of the data stored. This is the **only mode** that supports changing the data collection update
46 - frequency (`update_every`) **without losing** the previously stored metrics. For more details see [here](engine/).
42 +6. `dbengine`, (the default) data are in database files. The [Database Engine](engine/) works like a traditional
43 + database. There is some amount of RAM dedicated to data caching and indexing and the rest of the data reside
44 + compressed on disk. The number of history entries is not fixed in this case, but depends on the configured disk
45 + space and the effective compression ratio of the data stored. This is the **only mode** that supports changing the
46 + data collection update frequency (`update_every`) **without losing** the previously stored metrics. For more details
47 + see [here](engine/).
48
49 You can select the memory mode by editing `netdata.conf` and setting:
50
@@ -63,7 +64,7 @@ Embedded devices usually have very limited RAM resources available.
64 There are 2 settings for you to tweak:
65
66 1. `update every`, which controls the data collection frequency
66 -2. `history`, which controls the size of the database in RAM
67 +2. `history`, which controls the size of the database in RAM (except for `memory mode = dbengine`)
68
69 By default `update every = 1` and `history = 3600`. This gives you an hour of data with per second updates.
70
database/engine/README.md
+51
@@ -141,4 +141,55 @@ kern.maxfiles=65536
141
142 You can apply the settings by running `sysctl -p` or by rebooting.
143
144 +## Evaluation
145 +
146 +We have evaluated the performance of the `dbengine` API that the netdata daemon uses internally. This is **not** the
147 +web API of netdata. Our benchmarks ran on a **single** `dbengine` instance, multiple of which can be running in a
148 +netdata master server. We used a server with an AMD Ryzen Threadripper 2950X 16-Core Processor and 2 disk drives, a
149 +Seagate Constellation ES.3 2TB magnetic HDD and a SAMSUNG MZQLB960HAJR-00007 960GB NAND Flash SSD.
150 +
151 +For our workload, we defined 32 charts with 128 metrics each, giving us a total of 4096 metrics. We defined 1 worker
152 +thread per chart (32 threads) that generates new data points with a data generation interval of 1 second. The time axis
153 +of the time-series is emulated and accelerated so that the worker threads can generate as many data points as possible
154 +without delays.
155 +
156 +We also defined 32 worker threads that perform queries on random metrics with semi-random time ranges. The
157 +starting time of the query is randomly selected between the beginning of the time-series and the time of the latest data
158 +point. The ending time is randomly selected between 1 second and 1 hour after the starting time. The pseudo-random
159 +numbers are generated with a uniform distribution.
160 +
161 +The data are written to the database at the same time as they are read from it. This is a concurrent read/write mixed
162 +workload with a duration of 60 seconds. The faster `dbengine` runs, the bigger the dataset size becomes since more
163 +data points will be generated. We set a page cache size of 64MiB for the two disk-bound scenarios. This way, the dataset
164 +size of the metric data is much bigger than the RAM that is being used for caching so as to trigger I/O requests most
165 +of the time. In our final scenario, we set the page cache size to 16 GiB. That way, the dataset fits in the page cache
166 +so as to avoid all disk bottlenecks.
167 +
168 +The reported numbers are the following:
169 +
170 +| device | page cache | dataset | reads/sec | writes/sec |
171 +| :---: | :---: | ---: | ---: | ---: |
172 +| HDD | 64 MiB | 4.1 GiB | 813K | 18.0M |
173 +| SSD | 64 MiB | 9.8 GiB | 1.7M | 43.0M |
174 +| N/A | 16 GiB | 6.8 GiB |118.2M | 30.2M |
175 +
176 +where "reads/sec" is the number of metric data points being read from the database via its API per second and
177 +"writes/sec" is the number of metric data points being written to the database per second.
178 +
179 +Notice that the HDD numbers are pretty high and not much slower than the SSD numbers. This is thanks to the database
180 +engine design being optimized for rotating media. In the database engine disk I/O requests are:
181 +
182 +- asynchronous to mask the high I/O latency of HDDs.
183 +- mostly large to reduce the amount of HDD seeking time.
184 +- mostly sequential to reduce the amount of HDD seeking time.
185 +- compressed to reduce the amount of required throughput.
186 +
187 +As a result, the HDD is not thousands of times slower than the SSD, which is typical for other workloads.
188 +
189 +An interesting observation to make is that the CPU-bound run (16 GiB page cache) generates fewer data than the SSD run
190 +(6.8 GiB vs 9.8 GiB). The reason is that the 32 reader threads in the SSD scenario are more frequently blocked by I/O,
191 +and generate a read load of 1.7M/sec, whereas in the CPU-bound scenario the read load is 70 times higher at 118M/sec.
192 +Consequently, there is a significant degree of interference by the reader threads, that slow down the writer threads.
193 +This is also possible because the interference effects are greater than the SSD impact on data generation throughput.
194 +
195 [![analytics](https://www.google-analytics.com/collect?v=1&aip=1&t=pageview&_s=1&ds=github&dr=https%3A%2F%2Fgithub.com%2Fnetdata%2Fnetdata&dl=https%3A%2F%2Fmy-netdata.io%2Fgithub%2Fdatabase%2Fengine%2FREADME&_u=MAC~&cid=5792dfd7-8dc4-476b-af31-da2fdb9f93d2&tid=UA-64295674-3)](<>)
database/engine/rrdengine.c
-43
@@ -815,47 +815,6 @@ error_after_loop_init:
815 complete(&ctx->rrdengine_completion);
816 }
817
818 -
819 -#define NR_PAGES (256)
820 -static void basic_functional_test(struct rrdengine_instance *ctx)
821 -{
822 - int i, j, failed_validations;
823 - uuid_t uuid[NR_PAGES];
824 - void *buf;
825 - struct rrdeng_page_descr *handle[NR_PAGES];
826 - char uuid_str[UUID_STR_LEN];
827 - char backup[NR_PAGES][UUID_STR_LEN * 100]; /* backup storage for page data verification */
828 -
829 - for (i = 0 ; i < NR_PAGES ; ++i) {
830 - uuid_generate(uuid[i]);
831 - uuid_unparse_lower(uuid[i], uuid_str);
832 -// fprintf(stderr, "Generated uuid[%d]=%s\n", i, uuid_str);
833 - buf = rrdeng_create_page(ctx, &uuid[i], &handle[i]);
834 - /* Each page contains 10 times its own UUID stringified */
835 - for (j = 0 ; j < 100 ; ++j) {
836 - strcpy(buf + UUID_STR_LEN * j, uuid_str);
837 - strcpy(backup[i] + UUID_STR_LEN * j, uuid_str);
838 - }
839 - rrdeng_commit_page(ctx, handle[i], (Word_t)i);
840 - }
841 - fprintf(stderr, "\n********** CREATED %d METRIC PAGES ***********\n\n", NR_PAGES);
842 - failed_validations = 0;
843 - for (i = 0 ; i < NR_PAGES ; ++i) {
844 - buf = rrdeng_get_latest_page(ctx, &uuid[i], (void **)&handle[i]);
845 - if (NULL == buf) {
846 - ++failed_validations;
847 - fprintf(stderr, "Page %d was LOST.\n", i);
848 - }
849 - if (memcmp(backup[i], buf, UUID_STR_LEN * 100)) {
850 - ++failed_validations;
851 - fprintf(stderr, "Page %d data comparison with backup FAILED validation.\n", i);
852 - }
853 - rrdeng_put_page(ctx, handle[i]);
854 - }
855 - fprintf(stderr, "\n********** CORRECTLY VALIDATED %d/%d METRIC PAGES ***********\n\n",
856 - NR_PAGES - failed_validations, NR_PAGES);
857 -
858 -}
818 /* C entry point for development purposes
819 * make "LDFLAGS=-errdengine_main"
820 */
@@ -868,8 +827,6 @@ void rrdengine_main(void)
827 if (ret) {
828 exit(ret);
829 }
871 - basic_functional_test(ctx);
872 -
830 rrdeng_exit(ctx);
831 fprintf(stderr, "Hello world!");
832 exit(0);
database/engine/rrdenginelib.c
+2 -2
@@ -8,7 +8,7 @@ void print_page_cache_descr(struct rrdeng_page_descr *descr)
8 {
9 struct page_cache_descr *pg_cache_descr = descr->pg_cache_descr;
10 char uuid_str[UUID_STR_LEN];
11 - char str[BUFSIZE];
11 + char str[BUFSIZE + 1];
12 int pos = 0;
13
14 uuid_unparse_lower(*descr->id, uuid_str);
@@ -31,7 +31,7 @@ void print_page_cache_descr(struct rrdeng_page_descr *descr)
31 void print_page_descr(struct rrdeng_page_descr *descr)
32 {
33 char uuid_str[UUID_STR_LEN];
34 - char str[BUFSIZE];
34 + char str[BUFSIZE + 1];
35 int pos = 0;
36
37 uuid_unparse_lower(*descr->id, uuid_str);
database/rrd.c
+4
@@ -15,7 +15,11 @@ int rrd_delete_unupdated_dimensions = 0;
15
16 int default_rrd_update_every = UPDATE_EVERY;
17 int default_rrd_history_entries = RRD_DEFAULT_HISTORY_ENTRIES;
18 +#ifdef ENABLE_DBENGINE
19 +RRD_MEMORY_MODE default_rrd_memory_mode = RRD_MEMORY_MODE_DBENGINE;
20 +#else
21 RRD_MEMORY_MODE default_rrd_memory_mode = RRD_MEMORY_MODE_SAVE;
22 +#endif
23 int gap_when_lost_iterations_above = 1;
24
25