fix parallel execution of pg_ptrack_get_block()

This commit is contained in:
Anastasia
2018-01-16 14:16:50 +03:00
parent 7ee0803849
commit 1ff4dde4ad
5 changed files with 95 additions and 22 deletions
+13 -15
View File
@@ -69,15 +69,6 @@ static bool backup_in_progress = false;
/* Is pg_stop_backup() was sent */
static bool pg_stop_backup_is_sent = false;
typedef struct
{
const char *from_root;
const char *to_root;
parray *backup_files_list;
parray *prev_backup_filelist;
XLogRecPtr prev_backup_start_lsn;
} backup_files_args;
/*
* Backup routines
*/
@@ -656,6 +647,8 @@ do_backup_instance(void)
arg->backup_files_list = backup_files_list;
arg->prev_backup_filelist = prev_backup_filelist;
arg->prev_backup_start_lsn = prev_backup_start_lsn;
arg->thread_backup_conn = NULL;
arg->thread_cancel_conn = NULL;
backup_threads_args[i] = arg;
}
@@ -678,6 +671,8 @@ do_backup_instance(void)
for (i = 0; i < num_threads; i++)
{
pthread_join(backup_threads[i], NULL);
if (backup_threads_args[i]->thread_backup_conn != NULL)
pgut_disconnect(backup_threads_args[i]->thread_backup_conn);
pg_free(backup_threads_args[i]);
}
@@ -1956,7 +1951,8 @@ backup_files(void *arg)
if (file->is_datafile && !file->is_cfs)
{
/* backup block by block if datafile AND not compressed by cfs*/
if (!backup_data_file(arguments->from_root,
if (!backup_data_file(arguments,
arguments->from_root,
arguments->to_root, file,
arguments->prev_backup_start_lsn,
current.backup_mode))
@@ -2686,7 +2682,8 @@ get_last_ptrack_lsn(void)
}
char *
pg_ptrack_get_block(Oid dbOid,
pg_ptrack_get_block(backup_files_args *arguments,
Oid dbOid,
Oid tblsOid,
Oid relOid,
BlockNumber blknum,
@@ -2695,7 +2692,6 @@ pg_ptrack_get_block(Oid dbOid,
PGresult *res;
char *params[4];
char *result;
PGconn *tmp_conn = NULL;
params[0] = palloc(64);
params[1] = palloc(64);
@@ -2711,10 +2707,13 @@ pg_ptrack_get_block(Oid dbOid,
sprintf(params[2], "%i", relOid);
sprintf(params[3], "%u", blknum);
tmp_conn = pgut_connect(pgut_dbname);
if (arguments->thread_backup_conn == NULL)
arguments->thread_backup_conn = pgut_connect(pgut_dbname);
//elog(LOG, "db %i pg_ptrack_get_block(%i, %i, %u)",dbOid, tblsOid, relOid, blknum);
res = pgut_execute(tmp_conn, "SELECT pg_ptrack_get_block_2($1, $2, $3, $4)",
res = pgut_execute_parallel(arguments->thread_backup_conn,
arguments->thread_cancel_conn,
"SELECT pg_ptrack_get_block_2($1, $2, $3, $4)",
4, (const char **)params, true);
if (PQnfields(res) != 1)
@@ -2735,7 +2734,6 @@ pg_ptrack_get_block(Oid dbOid,
result_size);
PQclear(res);
pgut_disconnect(tmp_conn);
pfree(params[0]);
pfree(params[1]);
+7 -5
View File
@@ -222,7 +222,8 @@ read_page_from_file(pgFile *file, BlockNumber blknum,
* to the backup file.
*/
static void
backup_data_page(pgFile *file, XLogRecPtr prev_backup_start_lsn,
backup_data_page(backup_files_args *arguments,
pgFile *file, XLogRecPtr prev_backup_start_lsn,
BlockNumber blknum, BlockNumber nblocks,
FILE *in, FILE *out,
pg_crc32 *crc, int *n_skipped,
@@ -274,7 +275,7 @@ backup_data_page(pgFile *file, XLogRecPtr prev_backup_start_lsn,
free(page);
page = NULL;
page = (Page) pg_ptrack_get_block(file->dbOid, file->tblspcOid,
page = (Page) pg_ptrack_get_block(arguments, file->dbOid, file->tblspcOid,
file->relOid, absolute_blknum, &page_size);
if (page == NULL)
@@ -371,7 +372,8 @@ backup_data_page(pgFile *file, XLogRecPtr prev_backup_start_lsn,
* backup with special header.
*/
bool
backup_data_file(const char *from_root, const char *to_root,
backup_data_file(backup_files_args* arguments,
const char *from_root, const char *to_root,
pgFile *file, XLogRecPtr prev_backup_start_lsn,
BackupMode backup_mode)
{
@@ -453,7 +455,7 @@ backup_data_file(const char *from_root, const char *to_root,
{
for (blknum = 0; blknum < nblocks; blknum++)
{
backup_data_page(file, prev_backup_start_lsn, blknum,
backup_data_page(arguments, file, prev_backup_start_lsn, blknum,
nblocks, in, out, &(file->crc),
&n_blocks_skipped, backup_mode);
n_blocks_read++;
@@ -466,7 +468,7 @@ backup_data_file(const char *from_root, const char *to_root,
iter = datapagemap_iterate(&file->pagemap);
while (datapagemap_next(iter, &blknum))
{
backup_data_page(file, prev_backup_start_lsn, blknum,
backup_data_page(arguments, file, prev_backup_start_lsn, blknum,
nblocks, in, out, &(file->crc),
&n_blocks_skipped, backup_mode);
n_blocks_read++;
+15 -2
View File
@@ -242,6 +242,16 @@ typedef union DataPage
char data[BLCKSZ];
} DataPage;
typedef struct
{
const char *from_root;
const char *to_root;
parray *backup_files_list;
parray *prev_backup_filelist;
XLogRecPtr prev_backup_start_lsn;
PGconn *thread_backup_conn;
PGconn *thread_cancel_conn;
} backup_files_args;
/*
* return pointer that exceeds the length of prefix from character string.
@@ -323,7 +333,9 @@ extern const char *deparse_backup_mode(BackupMode mode);
extern void process_block_change(ForkNumber forknum, RelFileNode rnode,
BlockNumber blkno);
extern char *pg_ptrack_get_block(Oid dbOid, Oid tblsOid, Oid relOid, BlockNumber blknum,
extern char *pg_ptrack_get_block(backup_files_args *arguments,
Oid dbOid, Oid tblsOid, Oid relOid,
BlockNumber blknum,
size_t *result_size);
/* in restore.c */
extern int do_restore_or_validate(time_t target_backup_id,
@@ -428,7 +440,8 @@ extern int pgFileCompareLinked(const void *f1, const void *f2);
extern int pgFileCompareSize(const void *f1, const void *f2);
/* in data.c */
extern bool backup_data_file(const char *from_root, const char *to_root,
extern bool backup_data_file(backup_files_args* arguments,
const char *from_root, const char *to_root,
pgFile *file, XLogRecPtr prev_backup_start_lsn,
BackupMode backup_mode);
extern void restore_data_file(const char *from_root, const char *to_root,
+57
View File
@@ -1509,6 +1509,63 @@ pgut_set_port(const char *new_port)
port = new_port;
}
PGresult *
pgut_execute_parallel(PGconn* conn,
PGconn* cancel_conn, const char *query,
int nParams, const char **params,
bool text_result)
{
PGresult *res;
if (interrupted && !in_cleanup)
elog(ERROR, "interrupted");
/* write query to elog if verbose */
if (LOG_LEVEL_CONSOLE <= LOG || LOG_LEVEL_FILE <= LOG)
{
int i;
if (strchr(query, '\n'))
elog(LOG, "(query)\n%s", query);
else
elog(LOG, "(query) %s", query);
for (i = 0; i < nParams; i++)
elog(LOG, "\t(param:%d) = %s", i, params[i] ? params[i] : "(null)");
}
if (conn == NULL)
{
elog(ERROR, "not connected");
return NULL;
}
//on_before_exec(conn);
if (nParams == 0)
res = PQexec(conn, query);
else
res = PQexecParams(conn, query, nParams, NULL, params, NULL, NULL,
/*
* Specify zero to obtain results in text format,
* or one to obtain results in binary format.
*/
(text_result) ? 0 : 1);
//on_after_exec();
switch (PQresultStatus(res))
{
case PGRES_TUPLES_OK:
case PGRES_COMMAND_OK:
case PGRES_COPY_IN:
break;
default:
elog(ERROR, "query failed: %squery was: %s",
PQerrorMessage(conn), query);
break;
}
return res;
}
PGresult *
pgut_execute(PGconn* conn, const char *query, int nParams, const char **params,
bool text_result)
+3
View File
@@ -127,6 +127,9 @@ extern PGconn *pgut_connect_replication_extended(const char *pghost, const char
extern void pgut_disconnect(PGconn *conn);
extern PGresult *pgut_execute(PGconn* conn, const char *query, int nParams,
const char **params, bool text_result);
extern PGresult *pgut_execute_parallel(PGconn* conn, PGconn* cancel_conn,
const char *query, int nParams,
const char **params, bool text_result);
extern bool pgut_send(PGconn* conn, const char *query, int nParams, const char **params, int elevel);
extern void pgut_cancel(PGconn* conn);
extern int pgut_wait(int num, PGconn *connections[], struct timeval *timeout);