From 29b45fafbcf1f09a73d3fc0cb71272833a4b0dc3 Mon Sep 17 00:00:00 2001 From: Grigory Smolkin Date: Tue, 17 Mar 2020 12:51:16 +0300 Subject: [PATCH] [Issue #174] refactoring --- src/archive.c | 324 +++++++++++++++++++++++--------------------------- 1 file changed, 147 insertions(+), 177 deletions(-) diff --git a/src/archive.c b/src/archive.c index 029ab2ec..3edc8f80 100644 --- a/src/archive.c +++ b/src/archive.c @@ -13,15 +13,15 @@ #include "utils/thread.h" #include "instr_time.h" -static int push_wal_file_uncompressed(const char *wal_file_name, const char *pg_xlog_dir, +static int push_file_internal_uncompressed(const char *wal_file_name, const char *pg_xlog_dir, const char *archive_dir, bool overwrite, bool no_sync, int thread_num, uint32 archive_timeout); #ifdef HAVE_LIBZ -static int push_wal_file_gz(const char *wal_file_name, const char *pg_xlog_dir, +static int push_file_internal_gz(const char *wal_file_name, const char *pg_xlog_dir, const char *archive_dir, bool overwrite, bool no_sync, int compress_level, int thread_num, uint32 archive_timeout); #endif -static void *push_wal_segno(void *arg); +static void *push_files(void *arg); static void get_wal_file(const char *from_path, const char *to_path); #ifdef HAVE_LIBZ static const char *get_gz_error(gzFile gzf, int errnum); @@ -64,62 +64,17 @@ typedef struct typedef struct WALSegno { char name[MAXFNAMELEN]; - XLogSegNo segno; volatile pg_atomic_flag lock; } WALSegno; -static int push_wal_segno_internal(WALSegno *xlogfile, const char *archive_status_dir, +static int push_file(WALSegno *xlogfile, const char *archive_status_dir, const char *pg_xlog_dir, const char *archive_dir, - int thread_num, bool overwrite, bool no_sync, - uint32 archive_timeout, bool no_ready_rename, - bool is_compress, int compress_level); + bool overwrite, bool no_sync, uint32 archive_timeout, + bool no_ready_rename, bool is_compress, + int compress_level, int thread_num); -/* Look for files with '.ready' suffix in archive_status directory - * and pack such files into batch sized array. - */ -static parray * -setup_push_filelist(const char *archive_status_dir, int batch_size) -{ - int i; - parray *batch_files = parray_new(); - parray *status_files = parray_new(); - - /* get list of files from archive_status */ - dir_list_file(status_files, archive_status_dir, false, false, false, 0, FIO_DB_HOST); - parray_qsort(status_files, pgFileComparePath); - - for (i = 0; i < parray_num(status_files); i++) - { - int result = 0; - char filename[MAXFNAMELEN]; - char suffix[MAXFNAMELEN]; - WALSegno *xlogfile = NULL; - pgFile *file = (pgFile *) parray_get(status_files, i); - - result = sscanf(file->name, "%[^.]%s", (char *) &filename, (char *) &suffix); - - if (result != 2) - continue; - - if (strcmp(suffix, ".ready") != 0) - continue; - - xlogfile = palloc(sizeof(WALSegno)); - pg_atomic_init_flag(&xlogfile->lock); - - strncpy(xlogfile->name, filename, MAXFNAMELEN); - parray_append(batch_files, xlogfile); - - if (parray_num(batch_files) >= batch_size) - break; - } - - /* cleanup */ - parray_walk(status_files, pgFileFree); - parray_free(status_files); - - return batch_files; -} +static parray *setup_push_filelist(const char *archive_status_dir, + const char *first_file, int batch_size); /* * At this point, we already done one roundtrip to archive server @@ -193,83 +148,49 @@ do_archive_push(InstanceConfig *instance, char *wal_file_path, is_compress = true; #endif - /* Single-thread push - * There are two cases, when we don`t want to start multi-thread push: - * - number of threads in equal to 1; multithreading in remote mode isn`t cheap, - * establishing ssh connection can take 100-200ms, so running and terminating - * one thread using generic multithread approach can take - * almost as much time as copying itself. - * - file to push is not WAL file, but .history, .backup or .partial file. - * we do not apply compression to such files. - */ - if (num_threads > batch_size) num_threads = batch_size; + /* Setup filelist and locks */ + batch_files = setup_push_filelist(archive_status_dir, wal_file_name, batch_size); + elog(INFO, "PID [%d]: pg_probackup push file %s into archive, " "threads: %i, batch size: %i, compression: %s", my_pid, wal_file_name, num_threads, batch_size, is_compress ? "zlib" : "none"); - if (!IsXLogFileName(wal_file_name) || num_threads == 1) + /* Single-thread push + * We don`t want to start multi-thread push, if number of threads in equal to 1. + * Multithreading in remote mode isn`t cheap, + * establishing ssh connection can take 100-200ms, so running and terminating + * one thread using generic multithread approach can take + * almost as much time as copying itself. + */ + if (num_threads == 1) { - int rc; - /* do not apply compression to .backup, .history and .partial files */ - if (!IsXLogFileName(wal_file_name)) - is_compress = false; - INSTR_TIME_SET_CURRENT(start_time); - - if (is_compress) - rc = push_wal_file_gz(wal_file_name, pg_xlog_dir, - instance->arclog_path, overwrite, - no_sync, instance->compress_level, 1, - instance->archive_timeout); - else - rc = push_wal_file_uncompressed(wal_file_name, pg_xlog_dir, - instance->arclog_path, overwrite, - no_sync, 1, instance->archive_timeout); - - if (rc == 0) - n_total_pushed++; - else - n_total_skipped++; - - if (batch_size > 1) + for (i = 0; i < parray_num(batch_files); i++) { - batch_files = setup_push_filelist(archive_status_dir, batch_size); + int rc; + WALSegno *xlogfile = (WALSegno *) parray_get(batch_files, i); - for (i = 0; i < parray_num(batch_files); i++) - { - WALSegno *xlogfile = (WALSegno *) parray_get(batch_files, i); - - /* skip first wal file */ - if (strcmp(xlogfile->name, wal_file_name) == 0) - continue; - - rc = push_wal_segno_internal(xlogfile, archive_status_dir, - pg_xlog_dir, instance->arclog_path, - 1, overwrite, no_sync, - instance->archive_timeout, - no_ready_rename, is_compress, - instance->compress_level); - - if (rc == 0) - n_total_pushed++; - else - n_total_skipped++; - } + rc = push_file(xlogfile, archive_status_dir, + pg_xlog_dir, instance->arclog_path, + overwrite, no_sync, + instance->archive_timeout, + no_ready_rename || (strcmp(xlogfile->name, wal_file_name) == 0) ? true : false, + is_compress && IsXLogFileName(xlogfile->name) ? true : false, + instance->compress_level, 1); + if (rc == 0) + n_total_pushed++; + else + n_total_skipped++; } - INSTR_TIME_SET_CURRENT(end_time); - push_isok = true; goto push_done; } - /* setup filelist and locks */ - batch_files = setup_push_filelist(archive_status_dir, batch_size); - /* init thread args with its own segno */ threads = (pthread_t *) palloc(sizeof(pthread_t) * num_threads); threads_args = (archive_push_arg *) palloc(sizeof(archive_push_arg) * num_threads); @@ -305,7 +226,7 @@ do_archive_push(InstanceConfig *instance, char *wal_file_path, for (i = 0; i < num_threads; i++) { archive_push_arg *arg = &(threads_args[i]); - pthread_create(&threads[i], NULL, push_wal_segno, arg); + pthread_create(&threads[i], NULL, push_files, arg); } /* Wait threads */ @@ -322,16 +243,15 @@ do_archive_push(InstanceConfig *instance, char *wal_file_path, n_total_skipped += threads_args[i].n_skipped; } - fio_disconnect(); - INSTR_TIME_SET_CURRENT(end_time); - /* Note, that we are leaking memory here, * because pushing into archive is a very * time-sensetive operation, so we skip freeing stuff. */ push_done: + fio_disconnect(); /* calculate elapsed time */ + INSTR_TIME_SET_CURRENT(end_time); INSTR_TIME_SUBTRACT(end_time, start_time); push_time = INSTR_TIME_GET_DOUBLE(end_time); pretty_time_interval(push_time, pretty_time_str, 20); @@ -353,57 +273,35 @@ push_done: * Copy files from pg_wal to archive catalog with possible compression. */ static void * -push_wal_segno(void *arg) +push_files(void *arg) { int i; int rc; - char wal_file_dummy[MAXPGPATH]; - char wal_file_ready[MAXPGPATH]; - char wal_file_done[MAXPGPATH]; archive_push_arg *args = (archive_push_arg *) arg; for (i = 0; i < parray_num(args->files); i++) { + bool no_ready_rename = args->no_ready_rename; WALSegno *xlogfile = (WALSegno *) parray_get(args->files, i); if (!pg_atomic_test_set_flag(&xlogfile->lock)) continue; - join_path_components(wal_file_dummy, args->archive_status_dir, xlogfile->name); - snprintf(wal_file_ready, MAXPGPATH, "%s.%s", wal_file_dummy, "ready"); - snprintf(wal_file_done, MAXPGPATH, "%s.%s", wal_file_dummy, "done"); + /* Do not rename ready file of the first file, + * we do this to avoid flooding PostgreSQL log with + * warnings about ready file been missing. + */ + if (strcmp(args->first_filename, xlogfile->name) == 0) + no_ready_rename = true; - elog(LOG, "Thread [%d]: pushing file \"%s\"", args->thread_num, xlogfile->name); - - /* If compression is not required, then just copy it as is */ - if (!args->compress) - rc = push_wal_file_uncompressed(xlogfile->name, args->pg_xlog_dir, - args->archive_dir, args->overwrite, - args->no_sync, args->thread_num, - args->archive_timeout); -#ifdef HAVE_LIBZ - else - rc = push_wal_file_gz(xlogfile->name, args->pg_xlog_dir, - args->archive_dir, args->overwrite, args->no_sync, - args->compress_level, args->thread_num, - args->archive_timeout); -#endif - - /* take '--no-ready-rename' flag into account */ - if (!args->no_ready_rename && - /* don`t rename ready file for the first segment, - * postgres will complain via WARNING if we do that */ - (strcmp(args->first_filename, xlogfile->name) != 0)) - { - canonicalize_path(wal_file_ready); - canonicalize_path(wal_file_done); - /* It is ok to rename status file in archive_status directory */ - elog(VERBOSE, "Thread [%d]: Rename \"%s\" to \"%s\"", args->thread_num, - wal_file_ready, wal_file_done); - if (fio_rename(wal_file_ready, wal_file_done, FIO_DB_HOST) < 0) - elog(WARNING, "Thread [%d]: Cannot rename ready file \"%s\" to \"%s\": %s", - args->thread_num, wal_file_ready, wal_file_done, strerror(errno)); - } + rc = push_file(xlogfile, args->archive_status_dir, + args->pg_xlog_dir, args->archive_dir, + args->overwrite, args->no_sync, + args->archive_timeout, + no_ready_rename, + /* do not compress .backup, .partial and .history files */ + args->compress && IsXLogFileName(xlogfile->name) ? true : false, + args->compress_level, args->thread_num); if (rc == 0) args->n_pushed++; @@ -419,33 +317,29 @@ push_wal_segno(void *arg) } int -push_wal_segno_internal(WALSegno *xlogfile, const char *archive_status_dir, - const char *pg_xlog_dir, const char *archive_dir, - int thread_num, bool overwrite, bool no_sync, - uint32 archive_timeout, bool no_ready_rename, - bool is_compress, int compress_level) +push_file(WALSegno *xlogfile, const char *archive_status_dir, + const char *pg_xlog_dir, const char *archive_dir, + bool overwrite, bool no_sync, uint32 archive_timeout, + bool no_ready_rename, bool is_compress, + int compress_level, int thread_num) { int rc; char wal_file_dummy[MAXPGPATH]; - char wal_file_ready[MAXPGPATH]; - char wal_file_done[MAXPGPATH]; join_path_components(wal_file_dummy, archive_status_dir, xlogfile->name); - snprintf(wal_file_ready, MAXPGPATH, "%s.%s", wal_file_dummy, "ready"); - snprintf(wal_file_done, MAXPGPATH, "%s.%s", wal_file_dummy, "done"); elog(LOG, "Thread [%d]: pushing file \"%s\"", thread_num, xlogfile->name); /* If compression is not required, then just copy it as is */ if (!is_compress) - rc = push_wal_file_uncompressed(xlogfile->name, pg_xlog_dir, - archive_dir, overwrite, no_sync, - thread_num, archive_timeout); + rc = push_file_internal_uncompressed(xlogfile->name, pg_xlog_dir, + archive_dir, overwrite, no_sync, + thread_num, archive_timeout); #ifdef HAVE_LIBZ else - rc = push_wal_file_gz(xlogfile->name, pg_xlog_dir, archive_dir, - overwrite, no_sync, compress_level, - thread_num, archive_timeout); + rc = push_file_internal_gz(xlogfile->name, pg_xlog_dir, archive_dir, + overwrite, no_sync, compress_level, + thread_num, archive_timeout); #endif /* take '--no-ready-rename' flag into account */ @@ -453,11 +347,19 @@ push_wal_segno_internal(WALSegno *xlogfile, const char *archive_status_dir, /* don`t rename ready file for the first segment, * postgres will complain via WARNING if we do that */ { + char wal_file_ready[MAXPGPATH]; + char wal_file_done[MAXPGPATH]; + + snprintf(wal_file_ready, MAXPGPATH, "%s.%s", wal_file_dummy, "ready"); + snprintf(wal_file_done, MAXPGPATH, "%s.%s", wal_file_dummy, "done"); + canonicalize_path(wal_file_ready); canonicalize_path(wal_file_done); /* It is ok to rename status file in archive_status directory */ elog(VERBOSE, "Thread [%d]: Rename \"%s\" to \"%s\"", thread_num, wal_file_ready, wal_file_done); + + /* do not error out, if rename failed */ if (fio_rename(wal_file_ready, wal_file_done, FIO_DB_HOST) < 0) elog(WARNING, "Thread [%d]: Cannot rename ready file \"%s\" to \"%s\": %s", thread_num, wal_file_ready, wal_file_done, strerror(errno)); @@ -475,9 +377,9 @@ push_wal_segno_internal(WALSegno *xlogfile, const char *archive_status_dir, * has the same checksum */ int -push_wal_file_uncompressed(const char *wal_file_name, const char *pg_xlog_dir, - const char *archive_dir, bool overwrite, bool no_sync, - int thread_num, uint32 archive_timeout) +push_file_internal_uncompressed(const char *wal_file_name, const char *pg_xlog_dir, + const char *archive_dir, bool overwrite, bool no_sync, + int thread_num, uint32 archive_timeout) { FILE *in = NULL; int out = -1; @@ -616,6 +518,8 @@ part_opened: elog(LOG, "Thread [%d]: WAL file already exists in archive with the same " "checksum, skip pushing: \"%s\"", thread_num, from_fullpath); /* cleanup */ + fio_fclose(in); + fio_close(out); fio_unlink(to_fullpath_part, FIO_BACKUP_HOST); return 1; } @@ -713,9 +617,9 @@ part_opened: * has the same checksum */ int -push_wal_file_gz(const char *wal_file_name, const char *pg_xlog_dir, - const char *archive_dir, bool overwrite, bool no_sync, - int compress_level, int thread_num, uint32 archive_timeout) +push_file_internal_gz(const char *wal_file_name, const char *pg_xlog_dir, + const char *archive_dir, bool overwrite, bool no_sync, + int compress_level, int thread_num, uint32 archive_timeout) { FILE *in = NULL; gzFile out = NULL; @@ -756,7 +660,7 @@ push_wal_file_gz(const char *wal_file_name, const char *pg_xlog_dir, if (out == NULL) { if (errno != EEXIST) - elog(ERROR, "Thread [%d]: Failed to open temp WAL file \"%s\": %s", + elog(ERROR, "Thread [%d]: Cannot open temp WAL file \"%s\": %s", thread_num, to_fullpath_gz_part, strerror(errno)); /* Already existing destination temp file is not an error condition */ } @@ -861,6 +765,8 @@ part_opened: elog(LOG, "Thread [%d]: WAL file already exists in archive with the same " "checksum, skip pushing: \"%s\"", thread_num, from_fullpath); /* cleanup */ + fio_fclose(in); + fio_gzclose(out); fio_unlink(to_fullpath_gz_part, FIO_BACKUP_HOST); return 1; } @@ -874,6 +780,8 @@ part_opened: /* Overwriting is forbidden, * so we must unlink partial file and exit with error. */ + fio_fclose(in); + fio_gzclose(out); fio_unlink(to_fullpath_gz_part, FIO_BACKUP_HOST); elog(ERROR, "Thread [%d]: WAL file already exists in archive with " "different checksum: \"%s\"", thread_num, to_fullpath_gz); @@ -1203,3 +1111,65 @@ copy_file_attributes(const char *from_path, fio_location from_location, to_path, strerror(errno)); } } + +/* Look for files with '.ready' suffix in archive_status directory + * and pack such files into batch sized array. + */ +parray * +setup_push_filelist(const char *archive_status_dir, const char *first_file, + int batch_size) +{ + int i; + WALSegno *xlogfile = NULL; + parray *status_files = NULL; + parray *batch_files = parray_new(); + + /* guarantee that first filename is in batch list */ + xlogfile = palloc(sizeof(WALSegno)); + pg_atomic_init_flag(&xlogfile->lock); + strncpy(xlogfile->name, first_file, MAXFNAMELEN); + parray_append(batch_files, xlogfile); + + if (batch_size < 2) + return batch_files; + + /* get list of files from archive_status */ + status_files = parray_new(); + dir_list_file(status_files, archive_status_dir, false, false, false, 0, FIO_DB_HOST); + parray_qsort(status_files, pgFileComparePath); + + for (i = 0; i < parray_num(status_files); i++) + { + int result = 0; + char filename[MAXFNAMELEN]; + char suffix[MAXFNAMELEN]; + pgFile *file = (pgFile *) parray_get(status_files, i); + + /* first filename already in batch list */ + if (strcmp(file->name, first_file) == 0) + continue; + + result = sscanf(file->name, "%[^.]%s", (char *) &filename, (char *) &suffix); + + if (result != 2) + continue; + + if (strcmp(suffix, ".ready") != 0) + continue; + + xlogfile = palloc(sizeof(WALSegno)); + pg_atomic_init_flag(&xlogfile->lock); + + strncpy(xlogfile->name, filename, MAXFNAMELEN); + parray_append(batch_files, xlogfile); + + if (parray_num(batch_files) >= batch_size) + break; + } + + /* cleanup */ + parray_walk(status_files, pgFileFree); + parray_free(status_files); + + return batch_files; +}