@cryptotaxi247 / netdata-1 / commits / f7fa21b63

Improve data write (#19525)

* Use one worker thread to store data in datafile and journalfile * Check for write errors, attempt to retry but continue anyway if data cannot be written * On write error, do not add to open cache --------- Co-authored-by: Costa Tsaousis <costa@netdata.cloud>

Stelios Fragkakis committed Feb 5, 2025 at 09:11 UTC f7fa21b639f2cf1ff42963f9f76d86b4f63cdc4b
5 files changed +109 -126
src/database/engine/journalfile.c
+35 -37
@@ -1,34 +1,9 @@
1 // SPDX-License-Identifier: GPL-3.0-or-later
2 #include "rrdengine.h"
3
4 -static void after_extent_write_journalfile_v1_io(uv_fs_t* req)
5 -{
6 - worker_is_busy(RRDENG_FLUSH_TRANSACTION_BUFFER_CB);
7 -
8 - WAL *wal = req->data;
9 - struct generic_io_descriptor *io_descr = &wal->io_descr;
10 - struct rrdengine_instance *ctx = io_descr->ctx;
11 -
12 - netdata_log_debug(D_RRDENGINE, "%s: Journal block was written to disk.", __func__);
13 - if (req->result < 0) {
14 - ctx_io_error(ctx);
15 - netdata_log_error("DBENGINE: %s: uv_fs_write: %s", __func__, uv_strerror((int)req->result));
16 - } else {
17 - netdata_log_debug(D_RRDENGINE, "%s: Journal block was written to disk.", __func__);
18 - }
19 -
20 - uv_fs_req_cleanup(req);
21 - wal_release(wal);
22 -
23 - __atomic_sub_fetch(&ctx->atomic.extents_currently_being_flushed, 1, __ATOMIC_RELAXED);
24 -
25 - worker_is_idle();
26 -}
27 -
4 /* Careful to always call this before creating a new journal file */
29 -void journalfile_v1_extent_write(struct rrdengine_instance *ctx, struct rrdengine_datafile *datafile, WAL *wal, uv_loop_t *loop)
5 +void journalfile_v1_extent_write(struct rrdengine_instance *ctx, struct rrdengine_datafile *datafile, WAL *wal)
6 {
31 - int ret;
7 struct generic_io_descriptor *io_descr;
8 struct rrdengine_journalfile *journalfile = datafile->journalfile;
9
@@ -48,12 +23,35 @@ void journalfile_v1_extent_write(struct rrdengine_instance *ctx, struct rrdengin
23
24 io_descr->req.data = wal;
25 io_descr->data = journalfile;
51 - io_descr->completion = NULL;
26
27 io_descr->iov = uv_buf_init((void *)io_descr->buf, wal->buf_size);
54 - ret = uv_fs_write(loop, &io_descr->req, journalfile->file, &io_descr->iov, 1,
55 - (int64_t)io_descr->pos, after_extent_write_journalfile_v1_io);
56 - fatal_assert(-1 != ret);
28 +
29 + int retries = 10;
30 + int ret = -1;
31 + while (ret == -1 && --retries) {
32 + ret = uv_fs_write(NULL, &io_descr->req, journalfile->file, &io_descr->iov, 1, (int64_t)io_descr->pos, NULL);
33 + if (ret == -1)
34 + sleep_usec(300 * USEC_PER_MS);
35 + }
36 +
37 + bool jf_write_error = (ret == -1 || io_descr->req.result < 0);
38 +
39 + if (unlikely(jf_write_error)) {
40 + ctx_io_error(ctx);
41 + if (ret == -1)
42 + netdata_log_error(
43 + "DBENGINE: %s: uv_fs_write: failed to store metadata in journalfile %u, offset %ld",
44 + __func__,
45 + datafile->fileno,
46 + (int64_t)io_descr->pos);
47 + else
48 + netdata_log_error("DBENGINE: %s: uv_fs_write: %s", __func__, uv_strerror((int)io_descr->req.result));
49 + }
50 +
51 + uv_fs_req_cleanup(&io_descr->req);
52 + wal_release(wal);
53 + __atomic_sub_fetch(&ctx->atomic.extents_currently_being_flushed, 1, __ATOMIC_RELAXED);
54 + worker_is_idle();
55
56 ctx_current_disk_space_increase(ctx, wal->buf_size);
57 ctx_io_write_op_bytes(ctx, wal->buf_size);
@@ -534,13 +532,13 @@ int journalfile_destroy_unsafe(struct rrdengine_journalfile *journalfile, struct
532 journalfile_v2_generate_path(datafile, path_v2, sizeof(path));
533
534 if (journalfile->file) {
537 - ret = uv_fs_ftruncate(NULL, &req, journalfile->file, 0, NULL);
538 - if (ret < 0) {
539 - netdata_log_error("DBENGINE: uv_fs_ftruncate(%s): %s", path, uv_strerror(ret));
540 - ctx_fs_error(ctx);
541 - }
542 - uv_fs_req_cleanup(&req);
543 - (void) close_uv_file(datafile, journalfile->file);
535 + ret = uv_fs_ftruncate(NULL, &req, journalfile->file, 0, NULL);
536 + if (ret < 0) {
537 + netdata_log_error("DBENGINE: uv_fs_ftruncate(%s): %s", path, uv_strerror(ret));
538 + ctx_fs_error(ctx);
539 + }
540 + uv_fs_req_cleanup(&req);
541 + (void)close_uv_file(datafile, journalfile->file);
542 }
543
544 // This is the new journal v2 index file
src/database/engine/journalfile.h
+1 -1
@@ -141,7 +141,7 @@ struct wal;
141 void journalfile_v1_generate_path(struct rrdengine_datafile *datafile, char *str, size_t maxlen);
142 void journalfile_v2_generate_path(struct rrdengine_datafile *datafile, char *str, size_t maxlen);
143 struct rrdengine_journalfile *journalfile_alloc_and_init(struct rrdengine_datafile *datafile);
144 -void journalfile_v1_extent_write(struct rrdengine_instance *ctx, struct rrdengine_datafile *datafile, struct wal *wal, uv_loop_t *loop);
144 +void journalfile_v1_extent_write(struct rrdengine_instance *ctx, struct rrdengine_datafile *datafile, struct wal *wal);
145 int journalfile_close(struct rrdengine_journalfile *journalfile, struct rrdengine_datafile *datafile);
146 int journalfile_unlink(struct rrdengine_journalfile *journalfile);
147 int journalfile_destroy_unsafe(struct rrdengine_journalfile *journalfile, struct rrdengine_datafile *datafile);
src/database/engine/pagecache.c
-3
@@ -46,9 +46,6 @@ static void main_cache_flush_dirty_page_callback(PGC *cache __maybe_unused, PGC_
46 descr->page_length = pgd_disk_footprint(descr->pgd);
47
48 DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(base, descr, link.prev, link.next);
49 -
50 - // TODO: ask @stelfrag/@ktsaou about this.
51 - // internal_fatal(descr->page_length > RRDENG_BLOCK_SIZE, "DBENGINE: faulty page length calculation");
49 }
50
51 struct completion completion;
src/database/engine/rrdengine.c
+69 -78
@@ -620,19 +620,13 @@ static void journalfile_extent_build(struct rrdengine_instance *ctx, struct exte
620 crc32set(jf_trailer->checksum, crc);
621 }
622
623 -static void after_extent_flushed_to_open(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t* req __maybe_unused, int status __maybe_unused) {
624 - if(completion)
625 - completion_mark_complete(completion);
626 -
627 - if(ctx_is_available_for_queries(ctx))
628 - rrdeng_enq_cmd(ctx, RRDENG_OPCODE_DATABASE_ROTATE, NULL, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
629 -}
630 -
631 -static void *extent_flushed_to_open_tp_worker(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t *uv_work_req __maybe_unused) {
623 +static void extent_flushed_to_open_tp_worker(
624 + struct rrdengine_instance *ctx,
625 + struct extent_io_descriptor *xt_io_descr,
626 + bool have_error)
627 +{
628 worker_is_busy(UV_EVENT_DBENGINE_FLUSHED_TO_OPEN);
629
634 - uv_fs_t *uv_fs_request = data;
635 - struct extent_io_descriptor *xt_io_descr = uv_fs_request->data;
630 struct page_descr_with_data *descr;
631 struct rrdengine_datafile *datafile;
632 unsigned i;
@@ -644,7 +638,7 @@ static void *extent_flushed_to_open_tp_worker(struct rrdengine_instance *ctx __m
638 for (i = 0 ; i < xt_io_descr->descr_count ; ++i) {
639 descr = xt_io_descr->descr_array[i];
640
647 - if (likely(still_running))
641 + if (likely(still_running && !have_error))
642 pgc_open_add_hot_page(
643 (Word_t)ctx, descr->metric_id,
644 (time_t) (descr->start_time_ut / USEC_PER_SEC),
@@ -656,7 +650,6 @@ static void *extent_flushed_to_open_tp_worker(struct rrdengine_instance *ctx __m
650 page_descriptor_release(descr);
651 }
652
659 - uv_fs_req_cleanup(uv_fs_request);
653 posix_memfree(xt_io_descr->buf);
654 extent_io_descriptor_release(xt_io_descr);
655
@@ -668,39 +661,10 @@ static void *extent_flushed_to_open_tp_worker(struct rrdengine_instance *ctx __m
661 // we just finished a flushing on a datafile that is not the active one
662 rrdeng_enq_cmd(ctx, RRDENG_OPCODE_JOURNAL_INDEX, datafile, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
663
671 - return data;
664 + worker_is_idle();
665 }
666
667 // Main event loop callback
675 -static void after_extent_write_datafile_io(uv_fs_t *uv_fs_request) {
676 - worker_is_busy(RRDENG_OPCODE_MAX + RRDENG_OPCODE_EXTENT_WRITE);
677 -
678 - struct extent_io_descriptor *xt_io_descr = uv_fs_request->data;
679 - struct rrdengine_datafile *datafile = xt_io_descr->datafile;
680 - struct rrdengine_instance *ctx = datafile->ctx;
681 -
682 - if (uv_fs_request->result < 0) {
683 - ctx_io_error(ctx);
684 - netdata_log_error("DBENGINE: %s: uv_fs_write(): %s", __func__, uv_strerror((int)uv_fs_request->result));
685 - }
686 -
687 - journalfile_v1_extent_write(ctx, xt_io_descr->datafile, xt_io_descr->wal, &rrdeng_main.loop);
688 -
689 - spinlock_lock(&datafile->writers.spinlock);
690 - datafile->writers.running--;
691 - datafile->writers.flushed_to_open_running++;
692 - spinlock_unlock(&datafile->writers.spinlock);
693 -
694 - rrdeng_enq_cmd(xt_io_descr->ctx,
695 - RRDENG_OPCODE_FLUSHED_TO_OPEN,
696 - uv_fs_request,
697 - xt_io_descr->completion,
698 - STORAGE_PRIORITY_INTERNAL_DBENGINE,
699 - NULL,
700 - NULL);
701 -
702 - worker_is_idle();
703 -}
668
669 static bool datafile_is_full(struct rrdengine_instance *ctx, struct rrdengine_datafile *datafile) {
670 bool ret = false;
@@ -766,7 +730,8 @@ static struct rrdengine_datafile *get_datafile_to_write_extent(struct rrdengine_
730 /*
731 * Take a page list in a judy array and write them
732 */
769 -static struct extent_io_descriptor *datafile_extent_build(struct rrdengine_instance *ctx, struct page_descr_with_data *base, struct completion *completion) {
733 +static struct extent_io_descriptor *datafile_extent_build(struct rrdengine_instance *ctx, struct page_descr_with_data *base)
734 +{
735 int ret;
736 unsigned i, count, size_bytes, pos, real_io_size;
737 uint32_t uncompressed_payload_length, max_compressed_size, payload_offset;
@@ -790,9 +755,6 @@ static struct extent_io_descriptor *datafile_extent_build(struct rrdengine_insta
755 }
756
757 if (!count) {
793 - if (completion)
794 - completion_mark_complete(completion);
795 -
758 __atomic_sub_fetch(&ctx->atomic.extents_currently_being_flushed, 1, __ATOMIC_RELAXED);
759 return NULL;
760 }
@@ -885,7 +847,6 @@ static struct extent_io_descriptor *datafile_extent_build(struct rrdengine_insta
847
848 xt_io_descr->bytes = size_bytes;
849 xt_io_descr->uv_fs_request.data = xt_io_descr;
888 - xt_io_descr->completion = completion;
850
851 trailer = xt_io_descr->buf + size_bytes - sizeof(*trailer);
852 crc = crc32(0L, Z_NULL, 0);
@@ -902,27 +863,68 @@ static struct extent_io_descriptor *datafile_extent_build(struct rrdengine_insta
863 return xt_io_descr;
864 }
865
905 -static void after_extent_write(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t* uv_work_req __maybe_unused, int status __maybe_unused) {
906 - struct extent_io_descriptor *xt_io_descr = data;
907 -
908 - if(xt_io_descr) {
909 - int ret = uv_fs_write(&rrdeng_main.loop,
910 - &xt_io_descr->uv_fs_request,
911 - xt_io_descr->datafile->file,
912 - &xt_io_descr->iov,
913 - 1,
914 - (int64_t) xt_io_descr->pos,
915 - after_extent_write_datafile_io);
916 -
917 - fatal_assert(-1 != ret);
918 - }
866 +static void after_extent_write(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t* uv_work_req __maybe_unused, int status __maybe_unused)
867 +{
868 + ;
869 }
870
921 -static void *extent_write_tp_worker(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t *uv_work_req __maybe_unused) {
871 +static void *extent_write_tp_worker(
872 + struct rrdengine_instance *ctx,
873 + void *data,
874 + struct completion *completion __maybe_unused,
875 + uv_work_t *req __maybe_unused)
876 +{
877 worker_is_busy(UV_EVENT_DBENGINE_EXTENT_WRITE);
878 struct page_descr_with_data *base = data;
924 - struct extent_io_descriptor *xt_io_descr = datafile_extent_build(ctx, base, completion);
925 - return xt_io_descr;
879 + struct extent_io_descriptor *xt_io_descr = datafile_extent_build(ctx, base);
880 +
881 + if (!xt_io_descr)
882 + goto done;
883 +
884 + struct rrdengine_datafile *datafile = xt_io_descr->datafile;
885 +
886 + int retries = 10;
887 + int ret = -1;
888 + while (ret == -1 && --retries) {
889 + ret = uv_fs_write(NULL, &xt_io_descr->uv_fs_request, datafile->file, &xt_io_descr->iov, 1, (int64_t)xt_io_descr->pos, NULL);
890 + if (ret == -1)
891 + sleep_usec(300 * USEC_PER_MS);
892 + }
893 +
894 + bool df_write_error = (ret == -1 || xt_io_descr->uv_fs_request.result < 0);
895 +
896 + if (unlikely(df_write_error)) {
897 + ctx_io_error(ctx);
898 + if (ret == -1)
899 + netdata_log_error(
900 + "DBENGINE: %s: uv_fs_write: failed to store metrics in datafile %u, offset %ld",
901 + __func__,
902 + datafile->fileno,
903 + (int64_t)xt_io_descr->pos);
904 + else
905 + netdata_log_error(
906 + "DBENGINE: %s: uv_fs_write: %s", __func__, uv_strerror((int)xt_io_descr->uv_fs_request.result));
907 + }
908 +
909 + if (likely(!df_write_error)) {
910 + journalfile_v1_extent_write(ctx, datafile, xt_io_descr->wal);
911 + }
912 +
913 + spinlock_lock(&datafile->writers.spinlock);
914 + datafile->writers.running--;
915 + datafile->writers.flushed_to_open_running++;
916 + spinlock_unlock(&datafile->writers.spinlock);
917 +
918 + extent_flushed_to_open_tp_worker(ctx, xt_io_descr, df_write_error);
919 +
920 + if(ctx_is_available_for_queries(ctx))
921 + rrdeng_enq_cmd(ctx, RRDENG_OPCODE_DATABASE_ROTATE, NULL, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
922 +done:
923 + if(completion)
924 + completion_mark_complete(completion);
925 +
926 + worker_is_idle();
927 + return NULL;
928 }
929
930 static void after_database_rotate(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t* req __maybe_unused, int status __maybe_unused) {
@@ -1696,7 +1698,7 @@ static void retention_timer_cb(uv_timer_t *handle) {
1698 if (!localhost)
1699 return;
1700
1699 - worker_is_busy(RRDENG_TIMER_CB);
1701 + worker_is_busy(RRDENG_RETENTION_TIMER_CB);
1702 uv_stop(handle->loop);
1703 uv_update_time(handle->loop);
1704
@@ -1870,7 +1872,6 @@ void dbengine_event_loop(void* arg) {
1872 worker_register_job_name(RRDENG_OPCODE_QUERY, "query");
1873 worker_register_job_name(RRDENG_OPCODE_EXTENT_WRITE, "extent write");
1874 worker_register_job_name(RRDENG_OPCODE_EXTENT_READ, "extent read");
1873 - worker_register_job_name(RRDENG_OPCODE_FLUSHED_TO_OPEN, "flushed to open");
1875 worker_register_job_name(RRDENG_OPCODE_DATABASE_ROTATE, "db rotate");
1876 worker_register_job_name(RRDENG_OPCODE_JOURNAL_INDEX, "journal index");
1877 worker_register_job_name(RRDENG_OPCODE_FLUSH_MAIN, "flush init");
@@ -1884,7 +1885,6 @@ void dbengine_event_loop(void* arg) {
1885 worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_QUERY, "query cb");
1886 worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_EXTENT_WRITE, "extent write cb");
1887 worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_EXTENT_READ, "extent read cb");
1887 - worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_FLUSHED_TO_OPEN, "flushed to open cb");
1888 worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_DATABASE_ROTATE, "db rotate cb");
1889 worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_JOURNAL_INDEX, "journal index cb");
1890 worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_FLUSH_MAIN, "flush init cb");
@@ -1893,8 +1893,8 @@ void dbengine_event_loop(void* arg) {
1893 worker_register_job_name(RRDENG_OPCODE_MAX + RRDENG_OPCODE_CTX_QUIESCE, "ctx quiesce cb");
1894
1895 // special jobs
1896 + worker_register_job_name(RRDENG_RETENTION_TIMER_CB, "retention timer");
1897 worker_register_job_name(RRDENG_TIMER_CB, "timer");
1897 - worker_register_job_name(RRDENG_FLUSH_TRANSACTION_BUFFER_CB, "transaction buffer flush cb");
1898
1899 worker_register_job_custom_metric(RRDENG_OPCODES_WAITING, "opcodes waiting", "opcodes", WORKER_METRIC_ABSOLUTE);
1900 worker_register_job_custom_metric(RRDENG_WORKS_DISPATCHED, "works dispatched", "works", WORKER_METRIC_ABSOLUTE);
@@ -1938,15 +1938,6 @@ void dbengine_event_loop(void* arg) {
1938 break;
1939 }
1940
1941 - case RRDENG_OPCODE_FLUSHED_TO_OPEN: {
1942 - struct rrdengine_instance *ctx = cmd.ctx;
1943 - uv_fs_t *uv_fs_request = cmd.data;
1944 - struct extent_io_descriptor *xt_io_descr = uv_fs_request->data;
1945 - struct completion *completion = xt_io_descr->completion;
1946 - work_dispatch(ctx, uv_fs_request, completion, opcode, extent_flushed_to_open_tp_worker, after_extent_flushed_to_open);
1947 - break;
1948 - }
1949 -
1941 case RRDENG_OPCODE_FLUSH_MAIN: {
1942 if(rrdeng_main.flushes_running < pgc_max_flushers()) {
1943 rrdeng_main.flushes_running++;
src/database/engine/rrdengine.h
+4 -7
@@ -248,7 +248,6 @@ enum rrdeng_opcode {
248 RRDENG_OPCODE_QUERY,
249 RRDENG_OPCODE_EXTENT_WRITE,
250 RRDENG_OPCODE_EXTENT_READ,
251 - RRDENG_OPCODE_FLUSHED_TO_OPEN,
251 RRDENG_OPCODE_DATABASE_ROTATE,
252 RRDENG_OPCODE_JOURNAL_INDEX,
253 RRDENG_OPCODE_FLUSH_MAIN,
@@ -269,10 +268,10 @@ enum rrdeng_opcode {
268 // RRDENG_MAX_OPCODE + opcode : reserved for the callbacks of each opcode
269 // RRDENG_MAX_OPCODE + RRDENG_MAX_OPCODE : reserved for the timer
270 #define RRDENG_TIMER_CB (RRDENG_OPCODE_MAX + RRDENG_OPCODE_MAX)
272 -#define RRDENG_FLUSH_TRANSACTION_BUFFER_CB (RRDENG_TIMER_CB + 1)
273 -#define RRDENG_OPCODES_WAITING (RRDENG_TIMER_CB + 2)
274 -#define RRDENG_WORKS_DISPATCHED (RRDENG_TIMER_CB + 3)
275 -#define RRDENG_WORKS_EXECUTING (RRDENG_TIMER_CB + 4)
271 +#define RRDENG_OPCODES_WAITING (RRDENG_TIMER_CB + 1)
272 +#define RRDENG_WORKS_DISPATCHED (RRDENG_TIMER_CB + 2)
273 +#define RRDENG_WORKS_EXECUTING (RRDENG_TIMER_CB + 3)
274 +#define RRDENG_RETENTION_TIMER_CB (RRDENG_TIMER_CB + 4)
275
276 struct extent_io_data {
277 unsigned fileno;
@@ -291,7 +290,6 @@ struct extent_io_descriptor {
290 struct wal *wal;
291 uint64_t pos;
292 unsigned bytes;
294 - struct completion *completion;
293 unsigned descr_count;
294 struct page_descr_with_data *descr_array[MAX_PAGES_PER_EXTENT];
295 struct rrdengine_datafile *datafile;
@@ -306,7 +304,6 @@ struct generic_io_descriptor {
304 void *data;
305 uint64_t pos;
306 unsigned bytes;
309 - struct completion *completion;
307 };
308
309 typedef struct wal {