mirror of
https://github.com/postgrespro/pg_probackup.git
synced 2026-06-21 01:34:15 +02:00
[Issue #174] use "--archive-timeout" option to set timeout for partial file expire
This commit is contained in:
+72
-34
@@ -14,11 +14,11 @@
|
||||
|
||||
static int push_wal_file_internal(const char *wal_file_name, const char *pg_xlog_dir,
|
||||
const char *archive_dir, bool overwrite, bool no_sync,
|
||||
int thread_num);
|
||||
int thread_num, uint32 archive_timeout);
|
||||
#ifdef HAVE_LIBZ
|
||||
static int gz_push_wal_file_internal(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);
|
||||
int compress_level, int thread_num, uint32 archive_timeout);
|
||||
#endif
|
||||
static void *push_wal_segno(void *arg);
|
||||
static void get_wal_file(const char *from_path, const char *to_path);
|
||||
@@ -33,14 +33,17 @@ static void copy_file_attributes(const char *from_path,
|
||||
typedef struct
|
||||
{
|
||||
TimeLineID tli;
|
||||
XLogSegNo first_segno;
|
||||
uint32 xlog_seg_size;
|
||||
|
||||
const char *pg_xlog_dir;
|
||||
const char *archive_dir;
|
||||
const char *archive_status_dir;
|
||||
bool overwrite;
|
||||
bool compress;
|
||||
bool no_sync;
|
||||
bool no_ready_rename;
|
||||
uint32 archive_timeout;
|
||||
|
||||
CompressAlg compress_alg;
|
||||
int compress_level;
|
||||
@@ -73,8 +76,7 @@ typedef struct WALSegno
|
||||
* set archive_command to
|
||||
* 'pg_probackup archive-push -B /home/anastasia/backup --wal-file-name %f',
|
||||
* to move backups into arclog_path.
|
||||
* Where archlog_path is $BACKUP_PATH/wal/system_id.
|
||||
* Currently it just copies wal files to the new location.
|
||||
* Where archlog_path is $BACKUP_PATH/wal/instance_name
|
||||
*/
|
||||
int
|
||||
do_archive_push(InstanceConfig *instance, char *wal_file_path,
|
||||
@@ -84,6 +86,7 @@ do_archive_push(InstanceConfig *instance, char *wal_file_path,
|
||||
uint64 i;
|
||||
char current_dir[MAXPGPATH];
|
||||
char pg_xlog_dir[MAXPGPATH];
|
||||
char archive_status_dir[MAXPGPATH];
|
||||
uint64 system_id;
|
||||
bool is_compress = false;
|
||||
|
||||
@@ -125,6 +128,7 @@ do_archive_push(InstanceConfig *instance, char *wal_file_path,
|
||||
elog(ERROR, "Cannot use pglz for WAL compression");
|
||||
|
||||
join_path_components(pg_xlog_dir, current_dir, XLOGDIR);
|
||||
join_path_components(archive_status_dir, pg_xlog_dir, "archive_status");
|
||||
|
||||
/* Create 'archlog_path' directory. Do nothing if it already exists. */
|
||||
//fio_mkdir(instance->arclog_path, DIR_PERMISSION, FIO_BACKUP_HOST);
|
||||
@@ -161,10 +165,12 @@ do_archive_push(InstanceConfig *instance, char *wal_file_path,
|
||||
if (is_compress)
|
||||
gz_push_wal_file_internal(wal_file_name, pg_xlog_dir,
|
||||
instance->arclog_path, overwrite,
|
||||
no_sync, instance->compress_level, 1);
|
||||
no_sync, instance->compress_level, 1,
|
||||
instance->archive_timeout);
|
||||
else
|
||||
push_wal_file_internal(wal_file_name, pg_xlog_dir,
|
||||
instance->arclog_path, overwrite, no_sync, 1);
|
||||
instance->arclog_path, overwrite,
|
||||
no_sync, 1, instance->archive_timeout);
|
||||
|
||||
push_isok = true;
|
||||
total_pushed_files++;
|
||||
@@ -175,8 +181,8 @@ do_archive_push(InstanceConfig *instance, char *wal_file_path,
|
||||
if (IsXLogFileName(wal_file_name))
|
||||
GetXLogFromFileName(wal_file_name, &tli, &first_segno, instance->xlog_seg_size);
|
||||
|
||||
/* setup filelist and locks */
|
||||
files = parray_new();
|
||||
/* setup locks */
|
||||
for (i = first_segno; i < first_segno + batch_size; i++)
|
||||
{
|
||||
WALSegno *xlogfile = palloc(sizeof(WALSegno));
|
||||
@@ -198,13 +204,16 @@ do_archive_push(InstanceConfig *instance, char *wal_file_path,
|
||||
archive_push_arg *arg = &(threads_args[i]);
|
||||
|
||||
arg->tli = tli;
|
||||
arg->first_segno = first_segno;
|
||||
arg->xlog_seg_size = instance->xlog_seg_size;
|
||||
arg->archive_dir = instance->arclog_path;
|
||||
arg->pg_xlog_dir = pg_xlog_dir;
|
||||
arg->archive_status_dir = archive_status_dir;
|
||||
arg->overwrite = overwrite;
|
||||
arg->compress = is_compress;
|
||||
arg->no_sync = no_sync;
|
||||
arg->no_ready_rename = no_ready_rename;
|
||||
arg->archive_timeout = instance->archive_timeout;
|
||||
|
||||
arg->compress_alg = instance->compress_alg;
|
||||
arg->compress_level = instance->compress_level;
|
||||
@@ -234,8 +243,9 @@ do_archive_push(InstanceConfig *instance, char *wal_file_path,
|
||||
total_pushed_files += threads_args[i].n_pushed_files;
|
||||
}
|
||||
|
||||
/* Note, that we don`t do garbage collection here,
|
||||
* because pushing into archive is a very time-sensetive operation.
|
||||
/* Note, that we are leaking memory here,
|
||||
* because pushing into archive is a very
|
||||
* time-sensetive operation, so we skip freeing stuff.
|
||||
*/
|
||||
|
||||
push_done:
|
||||
@@ -254,7 +264,8 @@ push_done:
|
||||
/* ------------- INTERNAL FUNCTIONS ---------- */
|
||||
/*
|
||||
* Copy WAL segment from pgdata to archive catalog with possible compression.
|
||||
*
|
||||
* TODO: make it possible to be greedy here, i.e. pushing everything
|
||||
* that is ready to be pushed.
|
||||
*/
|
||||
static void *
|
||||
push_wal_segno(void *arg)
|
||||
@@ -262,15 +273,11 @@ push_wal_segno(void *arg)
|
||||
int i;
|
||||
int rc;
|
||||
char wal_filename[MAXPGPATH];
|
||||
|
||||
char archive_status_dir[MAXPGPATH];
|
||||
char wal_file_dummy[MAXPGPATH];
|
||||
char wal_file_ready[MAXPGPATH];
|
||||
char wal_file_done[MAXPGPATH];
|
||||
archive_push_arg *args = (archive_push_arg *) arg;
|
||||
|
||||
join_path_components(archive_status_dir, args->pg_xlog_dir, "archive_status");
|
||||
|
||||
for (i = 0; i < parray_num(args->files); i++)
|
||||
{
|
||||
WALSegno *xlogfile = (WALSegno *) parray_get(args->files, i);
|
||||
@@ -281,7 +288,7 @@ push_wal_segno(void *arg)
|
||||
/* At first we must construct WAL filename from segno, tli and xlog_seg_size */
|
||||
GetXLogFileName(wal_filename, args->tli, xlogfile->segno, args->xlog_seg_size);
|
||||
|
||||
join_path_components(wal_file_dummy, archive_status_dir, wal_filename);
|
||||
join_path_components(wal_file_dummy, args->archive_status_dir, wal_filename);
|
||||
snprintf(wal_file_ready, MAXPGPATH, "%s.%s", wal_file_dummy, "ready");
|
||||
snprintf(wal_file_done, MAXPGPATH, "%s.%s", wal_file_dummy, "done");
|
||||
|
||||
@@ -298,17 +305,24 @@ push_wal_segno(void *arg)
|
||||
if (!args->compress)
|
||||
rc = push_wal_file_internal(wal_filename, args->pg_xlog_dir,
|
||||
args->archive_dir, args->overwrite,
|
||||
args->no_sync, args->thread_num);
|
||||
args->no_sync, args->thread_num,
|
||||
args->archive_timeout);
|
||||
#ifdef HAVE_LIBZ
|
||||
else
|
||||
rc = gz_push_wal_file_internal(wal_filename, args->pg_xlog_dir,
|
||||
args->archive_dir, args->overwrite, args->no_sync,
|
||||
args->compress_level, args->thread_num);
|
||||
args->compress_level, args->thread_num,
|
||||
args->archive_timeout);
|
||||
#endif
|
||||
|
||||
/* take '--no-ready-rename' flag into account */
|
||||
if (!args->no_ready_rename && rc == 0)
|
||||
if (!args->no_ready_rename && rc == 0 &&
|
||||
/* don`t rename ready file for the first segment,
|
||||
* postgres will complain via WARNING if we do that */
|
||||
args->first_segno != xlogfile->segno)
|
||||
{
|
||||
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);
|
||||
@@ -318,7 +332,6 @@ push_wal_segno(void *arg)
|
||||
}
|
||||
|
||||
args->n_pushed_files++;
|
||||
|
||||
}
|
||||
|
||||
args->ret = 0;
|
||||
@@ -336,7 +349,7 @@ push_wal_segno(void *arg)
|
||||
int
|
||||
push_wal_file_internal(const char *wal_file_name, const char *pg_xlog_dir,
|
||||
const char *archive_dir, bool overwrite, bool no_sync,
|
||||
int thread_num)
|
||||
int thread_num, uint32 archive_timeout)
|
||||
{
|
||||
FILE *in = NULL;
|
||||
int out = -1;
|
||||
@@ -353,8 +366,10 @@ push_wal_file_internal(const char *wal_file_name, const char *pg_xlog_dir,
|
||||
|
||||
/* from path */
|
||||
join_path_components(from_fullpath, pg_xlog_dir, wal_file_name);
|
||||
canonicalize_path(from_fullpath);
|
||||
/* to path */
|
||||
join_path_components(to_fullpath, archive_dir, wal_file_name);
|
||||
canonicalize_path(to_fullpath);
|
||||
|
||||
/* Open source file for read */
|
||||
in = fio_fopen(from_fullpath, PG_BINARY_R, FIO_DB_HOST);
|
||||
@@ -377,11 +392,20 @@ push_wal_file_internal(const char *wal_file_name, const char *pg_xlog_dir,
|
||||
else
|
||||
goto part_opened;
|
||||
|
||||
/* Partial file already exists, it could have happened due to failed archive-push,
|
||||
* in this case partial file can be discarded, or due to concurrent archiving.
|
||||
/*
|
||||
* Partial file already exists, it could have happened due to:
|
||||
* 1. failed archive-push
|
||||
* 2. concurrent archiving
|
||||
*
|
||||
* For ARCHIVE_TIMEOUT period we will try to create partial file
|
||||
* and look for the size of already existing partial file, to
|
||||
* determine if it is changing or not.
|
||||
* If after ARCHIVE_TIMEOUT we still failed to create partial
|
||||
* file, we will make a decision about discarding
|
||||
* already existing partial file.
|
||||
*/
|
||||
/* TODO: use --archive-timeout */
|
||||
while (partial_try_count < PARTIAL_WAL_TIMER)
|
||||
|
||||
while (partial_try_count < archive_timeout)
|
||||
{
|
||||
if (fio_stat(to_fullpath_part, &st, false, FIO_BACKUP_HOST) < 0)
|
||||
{
|
||||
@@ -407,8 +431,9 @@ push_wal_file_internal(const char *wal_file_name, const char *pg_xlog_dir,
|
||||
/* first round */
|
||||
if (!partial_try_count)
|
||||
{
|
||||
elog(VERBOSE, "Thread [%d]: Temp WAL file already exists, waiting on it: \"%s\"",
|
||||
thread_num, to_fullpath_part);
|
||||
elog(VERBOSE, "Thread [%d]: Temp WAL file already exists, "
|
||||
"waiting on it %s seconds: \"%s\"",
|
||||
thread_num, archive_timeout, to_fullpath_part);
|
||||
partial_file_size = st.st_size;
|
||||
}
|
||||
|
||||
@@ -433,7 +458,7 @@ push_wal_file_internal(const char *wal_file_name, const char *pg_xlog_dir,
|
||||
{
|
||||
if (!partial_is_stale)
|
||||
elog(ERROR, "Thread [%d]: Failed to open temp WAL file \"%s\" in %i seconds",
|
||||
thread_num, to_fullpath_part, PARTIAL_WAL_TIMER);
|
||||
thread_num, to_fullpath_part, archive_timeout);
|
||||
|
||||
/* Partial segment is considered stale, so reuse it */
|
||||
elog(LOG, "Thread [%d]: Reusing stale temp WAL file \"%s\"",
|
||||
@@ -562,7 +587,7 @@ part_opened:
|
||||
int
|
||||
gz_push_wal_file_internal(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)
|
||||
int compress_level, int thread_num, uint32 archive_timeout)
|
||||
{
|
||||
FILE *in = NULL;
|
||||
gzFile out = NULL;
|
||||
@@ -582,8 +607,10 @@ gz_push_wal_file_internal(const char *wal_file_name, const char *pg_xlog_dir,
|
||||
|
||||
/* from path */
|
||||
join_path_components(from_fullpath, pg_xlog_dir, wal_file_name);
|
||||
canonicalize_path(from_fullpath);
|
||||
/* to path */
|
||||
join_path_components(to_fullpath, archive_dir, wal_file_name);
|
||||
canonicalize_path(to_fullpath);
|
||||
|
||||
/* destination file with .gz suffix */
|
||||
snprintf(to_fullpath_gz, sizeof(to_fullpath_gz), "%s.gz", to_fullpath);
|
||||
@@ -608,10 +635,20 @@ gz_push_wal_file_internal(const char *wal_file_name, const char *pg_xlog_dir,
|
||||
else
|
||||
goto part_opened;
|
||||
|
||||
/* Partial file already exists, it could have happened due to failed archive-push,
|
||||
* in this case partial file can be discarded, or due to concurrent archiving.
|
||||
/*
|
||||
* Partial file already exists, it could have happened due to:
|
||||
* 1. failed archive-push
|
||||
* 2. concurrent archiving
|
||||
*
|
||||
* For ARCHIVE_TIMEOUT period we will try to create partial file
|
||||
* and look for the size of already existing partial file, to
|
||||
* determine if it is changing or not.
|
||||
* If after ARCHIVE_TIMEOUT we still failed to create partial
|
||||
* file, we will make a decision about discarding
|
||||
* already existing partial file.
|
||||
*/
|
||||
while (partial_try_count < PARTIAL_WAL_TIMER)
|
||||
|
||||
while (partial_try_count < archive_timeout)
|
||||
{
|
||||
if (fio_stat(to_fullpath_gz_part, &st, false, FIO_BACKUP_HOST) < 0)
|
||||
{
|
||||
@@ -637,8 +674,9 @@ gz_push_wal_file_internal(const char *wal_file_name, const char *pg_xlog_dir,
|
||||
/* first round */
|
||||
if (!partial_try_count)
|
||||
{
|
||||
elog(VERBOSE, "Thread [%d]: Temp WAL file already exists, waiting on it: \"%s\"",
|
||||
thread_num, to_fullpath_gz_part);
|
||||
elog(VERBOSE, "Thread [%d]: Temp WAL file already exists, "
|
||||
"waiting on it %s seconds: \"%s\"",
|
||||
thread_num, archive_timeout, to_fullpath_gz_part);
|
||||
partial_file_size = st.st_size;
|
||||
}
|
||||
|
||||
@@ -663,7 +701,7 @@ gz_push_wal_file_internal(const char *wal_file_name, const char *pg_xlog_dir,
|
||||
{
|
||||
if (!partial_is_stale)
|
||||
elog(ERROR, "Thread [%d]: Failed to open temp WAL file \"%s\" in %i seconds",
|
||||
thread_num, to_fullpath_gz_part, PARTIAL_WAL_TIMER);
|
||||
thread_num, to_fullpath_gz_part, archive_timeout);
|
||||
|
||||
/* Partial segment is considered stale, so reuse it */
|
||||
elog(LOG, "Thread [%d]: Reusing stale temp WAL file \"%s\"",
|
||||
|
||||
@@ -67,7 +67,6 @@ extern const char *PROGRAM_EMAIL;
|
||||
#define DATABASE_MAP "database_map"
|
||||
|
||||
/* Timeout defaults */
|
||||
#define PARTIAL_WAL_TIMER 60
|
||||
#define ARCHIVE_TIMEOUT_DEFAULT 300
|
||||
#define REPLICA_TIMEOUT_DEFAULT 300
|
||||
|
||||
|
||||
Reference in New Issue
Block a user