mirror of
https://github.com/postgrespro/pg_probackup.git
synced 2026-06-21 01:34:15 +02:00
[Issue #174] refactoring
This commit is contained in:
+147
-177
@@ -13,15 +13,15 @@
|
|||||||
#include "utils/thread.h"
|
#include "utils/thread.h"
|
||||||
#include "instr_time.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,
|
const char *archive_dir, bool overwrite, bool no_sync,
|
||||||
int thread_num, uint32 archive_timeout);
|
int thread_num, uint32 archive_timeout);
|
||||||
#ifdef HAVE_LIBZ
|
#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,
|
const char *archive_dir, bool overwrite, bool no_sync,
|
||||||
int compress_level, int thread_num, uint32 archive_timeout);
|
int compress_level, int thread_num, uint32 archive_timeout);
|
||||||
#endif
|
#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);
|
static void get_wal_file(const char *from_path, const char *to_path);
|
||||||
#ifdef HAVE_LIBZ
|
#ifdef HAVE_LIBZ
|
||||||
static const char *get_gz_error(gzFile gzf, int errnum);
|
static const char *get_gz_error(gzFile gzf, int errnum);
|
||||||
@@ -64,62 +64,17 @@ typedef struct
|
|||||||
typedef struct WALSegno
|
typedef struct WALSegno
|
||||||
{
|
{
|
||||||
char name[MAXFNAMELEN];
|
char name[MAXFNAMELEN];
|
||||||
XLogSegNo segno;
|
|
||||||
volatile pg_atomic_flag lock;
|
volatile pg_atomic_flag lock;
|
||||||
} WALSegno;
|
} 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,
|
const char *pg_xlog_dir, const char *archive_dir,
|
||||||
int thread_num, bool overwrite, bool no_sync,
|
bool overwrite, bool no_sync, uint32 archive_timeout,
|
||||||
uint32 archive_timeout, bool no_ready_rename,
|
bool no_ready_rename, bool is_compress,
|
||||||
bool is_compress, int compress_level);
|
int compress_level, int thread_num);
|
||||||
|
|
||||||
/* Look for files with '.ready' suffix in archive_status directory
|
static parray *setup_push_filelist(const char *archive_status_dir,
|
||||||
* and pack such files into batch sized array.
|
const char *first_file, int batch_size);
|
||||||
*/
|
|
||||||
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;
|
|
||||||
}
|
|
||||||
|
|
||||||
/*
|
/*
|
||||||
* At this point, we already done one roundtrip to archive server
|
* 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;
|
is_compress = true;
|
||||||
#endif
|
#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)
|
if (num_threads > batch_size)
|
||||||
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, "
|
elog(INFO, "PID [%d]: pg_probackup push file %s into archive, "
|
||||||
"threads: %i, batch size: %i, compression: %s",
|
"threads: %i, batch size: %i, compression: %s",
|
||||||
my_pid, wal_file_name, num_threads,
|
my_pid, wal_file_name, num_threads,
|
||||||
batch_size, is_compress ? "zlib" : "none");
|
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);
|
INSTR_TIME_SET_CURRENT(start_time);
|
||||||
|
for (i = 0; i < parray_num(batch_files); i++)
|
||||||
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)
|
|
||||||
{
|
{
|
||||||
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++)
|
rc = push_file(xlogfile, archive_status_dir,
|
||||||
{
|
pg_xlog_dir, instance->arclog_path,
|
||||||
WALSegno *xlogfile = (WALSegno *) parray_get(batch_files, i);
|
overwrite, no_sync,
|
||||||
|
instance->archive_timeout,
|
||||||
/* skip first wal file */
|
no_ready_rename || (strcmp(xlogfile->name, wal_file_name) == 0) ? true : false,
|
||||||
if (strcmp(xlogfile->name, wal_file_name) == 0)
|
is_compress && IsXLogFileName(xlogfile->name) ? true : false,
|
||||||
continue;
|
instance->compress_level, 1);
|
||||||
|
if (rc == 0)
|
||||||
rc = push_wal_segno_internal(xlogfile, archive_status_dir,
|
n_total_pushed++;
|
||||||
pg_xlog_dir, instance->arclog_path,
|
else
|
||||||
1, overwrite, no_sync,
|
n_total_skipped++;
|
||||||
instance->archive_timeout,
|
|
||||||
no_ready_rename, is_compress,
|
|
||||||
instance->compress_level);
|
|
||||||
|
|
||||||
if (rc == 0)
|
|
||||||
n_total_pushed++;
|
|
||||||
else
|
|
||||||
n_total_skipped++;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
INSTR_TIME_SET_CURRENT(end_time);
|
|
||||||
|
|
||||||
push_isok = true;
|
push_isok = true;
|
||||||
goto push_done;
|
goto push_done;
|
||||||
}
|
}
|
||||||
|
|
||||||
/* setup filelist and locks */
|
|
||||||
batch_files = setup_push_filelist(archive_status_dir, batch_size);
|
|
||||||
|
|
||||||
/* init thread args with its own segno */
|
/* init thread args with its own segno */
|
||||||
threads = (pthread_t *) palloc(sizeof(pthread_t) * num_threads);
|
threads = (pthread_t *) palloc(sizeof(pthread_t) * num_threads);
|
||||||
threads_args = (archive_push_arg *) palloc(sizeof(archive_push_arg) * 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++)
|
for (i = 0; i < num_threads; i++)
|
||||||
{
|
{
|
||||||
archive_push_arg *arg = &(threads_args[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 */
|
/* Wait threads */
|
||||||
@@ -322,16 +243,15 @@ do_archive_push(InstanceConfig *instance, char *wal_file_path,
|
|||||||
n_total_skipped += threads_args[i].n_skipped;
|
n_total_skipped += threads_args[i].n_skipped;
|
||||||
}
|
}
|
||||||
|
|
||||||
fio_disconnect();
|
|
||||||
INSTR_TIME_SET_CURRENT(end_time);
|
|
||||||
|
|
||||||
/* Note, that we are leaking memory here,
|
/* Note, that we are leaking memory here,
|
||||||
* because pushing into archive is a very
|
* because pushing into archive is a very
|
||||||
* time-sensetive operation, so we skip freeing stuff.
|
* time-sensetive operation, so we skip freeing stuff.
|
||||||
*/
|
*/
|
||||||
|
|
||||||
push_done:
|
push_done:
|
||||||
|
fio_disconnect();
|
||||||
/* calculate elapsed time */
|
/* calculate elapsed time */
|
||||||
|
INSTR_TIME_SET_CURRENT(end_time);
|
||||||
INSTR_TIME_SUBTRACT(end_time, start_time);
|
INSTR_TIME_SUBTRACT(end_time, start_time);
|
||||||
push_time = INSTR_TIME_GET_DOUBLE(end_time);
|
push_time = INSTR_TIME_GET_DOUBLE(end_time);
|
||||||
pretty_time_interval(push_time, pretty_time_str, 20);
|
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.
|
* Copy files from pg_wal to archive catalog with possible compression.
|
||||||
*/
|
*/
|
||||||
static void *
|
static void *
|
||||||
push_wal_segno(void *arg)
|
push_files(void *arg)
|
||||||
{
|
{
|
||||||
int i;
|
int i;
|
||||||
int rc;
|
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;
|
archive_push_arg *args = (archive_push_arg *) arg;
|
||||||
|
|
||||||
for (i = 0; i < parray_num(args->files); i++)
|
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);
|
WALSegno *xlogfile = (WALSegno *) parray_get(args->files, i);
|
||||||
|
|
||||||
if (!pg_atomic_test_set_flag(&xlogfile->lock))
|
if (!pg_atomic_test_set_flag(&xlogfile->lock))
|
||||||
continue;
|
continue;
|
||||||
|
|
||||||
join_path_components(wal_file_dummy, args->archive_status_dir, xlogfile->name);
|
/* Do not rename ready file of the first file,
|
||||||
snprintf(wal_file_ready, MAXPGPATH, "%s.%s", wal_file_dummy, "ready");
|
* we do this to avoid flooding PostgreSQL log with
|
||||||
snprintf(wal_file_done, MAXPGPATH, "%s.%s", wal_file_dummy, "done");
|
* 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);
|
rc = push_file(xlogfile, args->archive_status_dir,
|
||||||
|
args->pg_xlog_dir, args->archive_dir,
|
||||||
/* If compression is not required, then just copy it as is */
|
args->overwrite, args->no_sync,
|
||||||
if (!args->compress)
|
args->archive_timeout,
|
||||||
rc = push_wal_file_uncompressed(xlogfile->name, args->pg_xlog_dir,
|
no_ready_rename,
|
||||||
args->archive_dir, args->overwrite,
|
/* do not compress .backup, .partial and .history files */
|
||||||
args->no_sync, args->thread_num,
|
args->compress && IsXLogFileName(xlogfile->name) ? true : false,
|
||||||
args->archive_timeout);
|
args->compress_level, args->thread_num);
|
||||||
#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));
|
|
||||||
}
|
|
||||||
|
|
||||||
if (rc == 0)
|
if (rc == 0)
|
||||||
args->n_pushed++;
|
args->n_pushed++;
|
||||||
@@ -419,33 +317,29 @@ push_wal_segno(void *arg)
|
|||||||
}
|
}
|
||||||
|
|
||||||
int
|
int
|
||||||
push_wal_segno_internal(WALSegno *xlogfile, const char *archive_status_dir,
|
push_file(WALSegno *xlogfile, const char *archive_status_dir,
|
||||||
const char *pg_xlog_dir, const char *archive_dir,
|
const char *pg_xlog_dir, const char *archive_dir,
|
||||||
int thread_num, bool overwrite, bool no_sync,
|
bool overwrite, bool no_sync, uint32 archive_timeout,
|
||||||
uint32 archive_timeout, bool no_ready_rename,
|
bool no_ready_rename, bool is_compress,
|
||||||
bool is_compress, int compress_level)
|
int compress_level, int thread_num)
|
||||||
{
|
{
|
||||||
int rc;
|
int rc;
|
||||||
char wal_file_dummy[MAXPGPATH];
|
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);
|
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);
|
elog(LOG, "Thread [%d]: pushing file \"%s\"", thread_num, xlogfile->name);
|
||||||
|
|
||||||
/* If compression is not required, then just copy it as is */
|
/* If compression is not required, then just copy it as is */
|
||||||
if (!is_compress)
|
if (!is_compress)
|
||||||
rc = push_wal_file_uncompressed(xlogfile->name, pg_xlog_dir,
|
rc = push_file_internal_uncompressed(xlogfile->name, pg_xlog_dir,
|
||||||
archive_dir, overwrite, no_sync,
|
archive_dir, overwrite, no_sync,
|
||||||
thread_num, archive_timeout);
|
thread_num, archive_timeout);
|
||||||
#ifdef HAVE_LIBZ
|
#ifdef HAVE_LIBZ
|
||||||
else
|
else
|
||||||
rc = push_wal_file_gz(xlogfile->name, pg_xlog_dir, archive_dir,
|
rc = push_file_internal_gz(xlogfile->name, pg_xlog_dir, archive_dir,
|
||||||
overwrite, no_sync, compress_level,
|
overwrite, no_sync, compress_level,
|
||||||
thread_num, archive_timeout);
|
thread_num, archive_timeout);
|
||||||
#endif
|
#endif
|
||||||
|
|
||||||
/* take '--no-ready-rename' flag into account */
|
/* 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,
|
/* don`t rename ready file for the first segment,
|
||||||
* postgres will complain via WARNING if we do that */
|
* 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_ready);
|
||||||
canonicalize_path(wal_file_done);
|
canonicalize_path(wal_file_done);
|
||||||
/* It is ok to rename status file in archive_status directory */
|
/* It is ok to rename status file in archive_status directory */
|
||||||
elog(VERBOSE, "Thread [%d]: Rename \"%s\" to \"%s\"", thread_num,
|
elog(VERBOSE, "Thread [%d]: Rename \"%s\" to \"%s\"", thread_num,
|
||||||
wal_file_ready, wal_file_done);
|
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)
|
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",
|
elog(WARNING, "Thread [%d]: Cannot rename ready file \"%s\" to \"%s\": %s",
|
||||||
thread_num, wal_file_ready, wal_file_done, strerror(errno));
|
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
|
* has the same checksum
|
||||||
*/
|
*/
|
||||||
int
|
int
|
||||||
push_wal_file_uncompressed(const char *wal_file_name, const char *pg_xlog_dir,
|
push_file_internal_uncompressed(const char *wal_file_name, const char *pg_xlog_dir,
|
||||||
const char *archive_dir, bool overwrite, bool no_sync,
|
const char *archive_dir, bool overwrite, bool no_sync,
|
||||||
int thread_num, uint32 archive_timeout)
|
int thread_num, uint32 archive_timeout)
|
||||||
{
|
{
|
||||||
FILE *in = NULL;
|
FILE *in = NULL;
|
||||||
int out = -1;
|
int out = -1;
|
||||||
@@ -616,6 +518,8 @@ part_opened:
|
|||||||
elog(LOG, "Thread [%d]: WAL file already exists in archive with the same "
|
elog(LOG, "Thread [%d]: WAL file already exists in archive with the same "
|
||||||
"checksum, skip pushing: \"%s\"", thread_num, from_fullpath);
|
"checksum, skip pushing: \"%s\"", thread_num, from_fullpath);
|
||||||
/* cleanup */
|
/* cleanup */
|
||||||
|
fio_fclose(in);
|
||||||
|
fio_close(out);
|
||||||
fio_unlink(to_fullpath_part, FIO_BACKUP_HOST);
|
fio_unlink(to_fullpath_part, FIO_BACKUP_HOST);
|
||||||
return 1;
|
return 1;
|
||||||
}
|
}
|
||||||
@@ -713,9 +617,9 @@ part_opened:
|
|||||||
* has the same checksum
|
* has the same checksum
|
||||||
*/
|
*/
|
||||||
int
|
int
|
||||||
push_wal_file_gz(const char *wal_file_name, const char *pg_xlog_dir,
|
push_file_internal_gz(const char *wal_file_name, const char *pg_xlog_dir,
|
||||||
const char *archive_dir, bool overwrite, bool no_sync,
|
const char *archive_dir, bool overwrite, bool no_sync,
|
||||||
int compress_level, int thread_num, uint32 archive_timeout)
|
int compress_level, int thread_num, uint32 archive_timeout)
|
||||||
{
|
{
|
||||||
FILE *in = NULL;
|
FILE *in = NULL;
|
||||||
gzFile out = 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 (out == NULL)
|
||||||
{
|
{
|
||||||
if (errno != EEXIST)
|
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));
|
thread_num, to_fullpath_gz_part, strerror(errno));
|
||||||
/* Already existing destination temp file is not an error condition */
|
/* 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 "
|
elog(LOG, "Thread [%d]: WAL file already exists in archive with the same "
|
||||||
"checksum, skip pushing: \"%s\"", thread_num, from_fullpath);
|
"checksum, skip pushing: \"%s\"", thread_num, from_fullpath);
|
||||||
/* cleanup */
|
/* cleanup */
|
||||||
|
fio_fclose(in);
|
||||||
|
fio_gzclose(out);
|
||||||
fio_unlink(to_fullpath_gz_part, FIO_BACKUP_HOST);
|
fio_unlink(to_fullpath_gz_part, FIO_BACKUP_HOST);
|
||||||
return 1;
|
return 1;
|
||||||
}
|
}
|
||||||
@@ -874,6 +780,8 @@ part_opened:
|
|||||||
/* Overwriting is forbidden,
|
/* Overwriting is forbidden,
|
||||||
* so we must unlink partial file and exit with error.
|
* 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);
|
fio_unlink(to_fullpath_gz_part, FIO_BACKUP_HOST);
|
||||||
elog(ERROR, "Thread [%d]: WAL file already exists in archive with "
|
elog(ERROR, "Thread [%d]: WAL file already exists in archive with "
|
||||||
"different checksum: \"%s\"", thread_num, to_fullpath_gz);
|
"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));
|
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;
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user