/* an entry in the last-start-times shared hash table */ typedefstruct LauncherLastStartTimesEntry
{
Oid subid; /* OID of logrep subscription (hash key) */
TimestampTz last_start_time; /* last time its apply worker was started */
} LauncherLastStartTimesEntry;
sub = (Subscription *) palloc0(sizeof(Subscription));
sub->oid = subform->oid;
sub->dbid = subform->subdbid;
sub->owner = subform->subowner;
sub->enabled = subform->subenabled;
sub->name = pstrdup(NameStr(subform->subname)); /* We don't fill fields we are not interested in. */
res = lappend(res, sub);
MemoryContextSwitchTo(oldcxt);
}
for (;;)
{
BgwHandleStatus status;
pid_t pid; int rc;
CHECK_FOR_INTERRUPTS();
LWLockAcquire(LogicalRepWorkerLock, LW_SHARED);
/* Worker either died or has started. Return false if died. */ if (!worker->in_use || worker->proc)
{
result = worker->in_use;
LWLockRelease(LogicalRepWorkerLock); break;
}
LWLockRelease(LogicalRepWorkerLock);
/* Check if worker has died before attaching, and clean up after it. */
status = GetBackgroundWorkerPid(handle, &pid);
if (status == BGWH_STOPPED)
{
LWLockAcquire(LogicalRepWorkerLock, LW_EXCLUSIVE); /* Ensure that this was indeed the worker we waited for. */ if (generation == worker->generation)
logicalrep_worker_cleanup(worker);
LWLockRelease(LogicalRepWorkerLock); break; /* result is already false */
}
/* Search for attached worker for a given subscription id. */ for (i = 0; i < max_logical_replication_workers; i++)
{
LogicalRepWorker *w = &LogicalRepCtx->workers[i];
/* Skip parallel apply workers. */ if (isParallelApplyWorker(w)) continue;
/* *Similartologicalrep_worker_find(),butreturnsalistofallworkersfor *thesubscription,insteadofjustone.
*/
List *
logicalrep_workers_find(Oid subid, bool only_running, bool acquire_lock)
{ int i;
List *res = NIL;
if (acquire_lock)
LWLockAcquire(LogicalRepWorkerLock, LW_SHARED);
Assert(LWLockHeldByMe(LogicalRepWorkerLock));
/* Search for attached worker for a given subscription id. */ for (i = 0; i < max_logical_replication_workers; i++)
{
LogicalRepWorker *w = &LogicalRepCtx->workers[i];
if (w->in_use && w->subid == subid && (!only_running || w->proc))
res = lappend(res, w);
}
if (acquire_lock)
LWLockRelease(LogicalRepWorkerLock);
return res;
}
/* *Startnewlogicalreplicationbackgroundworker,ifpossible. * *Returnstrueonsuccess,falseonfailure.
*/ bool
logicalrep_worker_launch(LogicalRepWorkerType wtype,
Oid dbid, Oid subid, constchar *subname, Oid userid,
Oid relid, dsm_handle subworker_dsm)
{
BackgroundWorker bgw;
BackgroundWorkerHandle *bgw_handle;
uint16 generation; int i; int slot = 0;
LogicalRepWorker *worker = NULL; int nsyncworkers; int nparallelapplyworkers;
TimestampTz now; bool is_tablesync_worker = (wtype == WORKERTYPE_TABLESYNC); bool is_parallel_apply_worker = (wtype == WORKERTYPE_PARALLEL_APPLY);
ereport(DEBUG1,
(errmsg_internal("starting logical replication worker for subscription \"%s\"",
subname)));
/* Report this after the initial starting message for consistency. */ if (max_active_replication_origins == 0)
ereport(ERROR,
(errcode(ERRCODE_CONFIGURATION_LIMIT_EXCEEDED),
errmsg("cannot start logical replication workers when \"max_active_replication_origins\" is 0")));
for (i = 0; i < max_logical_replication_workers; i++)
{
LogicalRepWorker *w = &LogicalRepCtx->workers[i];
/* *Iftheworkerwasmarkedinusebutdidn'tmanagetoattachin *time,cleanitup.
*/ if (w->in_use && !w->proc &&
TimestampDifferenceExceeds(w->launch_time, now,
wal_receiver_timeout))
{
elog(WARNING, "logical replication worker for subscription %u took too long to start; canceled",
w->subid);
if (!RegisterDynamicBackgroundWorker(&bgw, &bgw_handle))
{ /* Failed to start worker, so clean up the worker slot. */
LWLockAcquire(LogicalRepWorkerLock, LW_EXCLUSIVE);
Assert(generation == worker->generation);
logicalrep_worker_cleanup(worker);
LWLockRelease(LogicalRepWorkerLock);
ereport(WARNING,
(errcode(ERRCODE_CONFIGURATION_LIMIT_EXCEEDED),
errmsg("out of background worker slots"),
errhint("You might need to increase \"%s\".", "max_worker_processes"))); returnfalse;
}
/* Now wait until it attaches. */ return WaitForReplicationWorkerAttach(worker, generation, bgw_handle);
}
/* *Ifwefoundaworkerbutitdoesnothaveprocsetthenitisstill *startingup;waitforittofinishstartingandthenkillit.
*/ while (worker->in_use && !worker->proc)
{ int rc;
LWLockRelease(LogicalRepWorkerLock);
/* Wait a bit --- we don't expect to have to wait long. */
rc = WaitLatch(MyLatch,
WL_LATCH_SET | WL_TIMEOUT | WL_EXIT_ON_PM_DEATH, 10L, WAIT_EVENT_BGWORKER_STARTUP);
if (rc & WL_LATCH_SET)
{
ResetLatch(MyLatch);
CHECK_FOR_INTERRUPTS();
}
/* Worker has assigned proc, so it has started. */ if (worker->proc) break;
}
/* Now terminate the worker ... */
kill(worker->proc->pid, signo);
/* ... and wait for it to die. */ for (;;)
{ int rc;
/* is it gone? */ if (!worker->proc || worker->generation != generation) break;
LWLockRelease(LogicalRepWorkerLock);
/* Wait a bit --- we don't expect to have to wait long. */
rc = WaitLatch(MyLatch,
WL_LATCH_SET | WL_TIMEOUT | WL_EXIT_ON_PM_DEATH, 10L, WAIT_EVENT_BGWORKER_SHUTDOWN);
if (rc & WL_LATCH_SET)
{
ResetLatch(MyLatch);
CHECK_FOR_INTERRUPTS();
}
/* *Cleanupfunction. * *Calledonlogicalreplicationworkerexit.
*/ staticvoid
logicalrep_worker_onexit(int code, Datum arg)
{ /* Disconnect gracefully from the remote side. */ if (LogRepWorkerWalRcvConn)
walrcv_disconnect(LogRepWorkerWalRcvConn);
logicalrep_worker_detach();
/* Cleanup fileset used for streaming transactions. */ if (MyLogicalRepWorker->stream_fileset != NULL)
FileSetDeleteAll(MyLogicalRepWorker->stream_fileset);
/* *Countthenumberofregistered(notnecessarilyrunning)syncworkers *forasubscription.
*/ int
logicalrep_sync_worker_count(Oid subid)
{ int i; int res = 0;
Assert(LWLockHeldByMe(LogicalRepWorkerLock));
/* Search for attached worker for a given subscription id. */ for (i = 0; i < max_logical_replication_workers; i++)
{
LogicalRepWorker *w = &LogicalRepCtx->workers[i];
if (isTablesyncWorker(w) && w->subid == subid)
res++;
}
return res;
}
/* *Countthenumberofregistered(butnotnecessarilyrunning)parallelapply *workersforasubscription.
*/ staticint
logicalrep_pa_worker_count(Oid subid)
{ int i; int res = 0;
Assert(LWLockHeldByMe(LogicalRepWorkerLock));
/* *Scanallattachedparallelapplyworkers,onlycountingthosewhich *havethegivensubscriptionid.
*/ for (i = 0; i < max_logical_replication_workers; i++)
{
LogicalRepWorker *w = &LogicalRepCtx->workers[i];
if (isParallelApplyWorker(w) && w->subid == subid)
res++;
}
/* Enter main loop */ for (;;)
{ int rc;
List *sublist;
ListCell *lc;
MemoryContext subctx;
MemoryContext oldctx; long wait_time = DEFAULT_NAPTIME_PER_CYCLE;
CHECK_FOR_INTERRUPTS();
/* Use temporary context to avoid leaking memory across cycles. */
subctx = AllocSetContextCreate(TopMemoryContext, "Logical Replication Launcher sublist",
ALLOCSET_DEFAULT_SIZES);
oldctx = MemoryContextSwitchTo(subctx);
/* Start any missing workers for enabled subscriptions. */
sublist = get_subscription_list();
foreach(lc, sublist)
{
Subscription *sub = (Subscription *) lfirst(lc);
LogicalRepWorker *w;
TimestampTz last_start;
TimestampTz now; long elapsed;
if (!sub->enabled) continue;
LWLockAcquire(LogicalRepWorkerLock, LW_SHARED);
w = logicalrep_worker_find(sub->oid, InvalidOid, false);
LWLockRelease(LogicalRepWorkerLock);
if (w != NULL) continue; /* worker is running already */
/* *Returnsstateofthesubscriptions.
*/
Datum
pg_stat_get_subscription(PG_FUNCTION_ARGS)
{ #define PG_STAT_GET_SUBSCRIPTION_COLS 10
Oid subid = PG_ARGISNULL(0) ? InvalidOid : PG_GETARG_OID(0); int i;
ReturnSetInfo *rsinfo = (ReturnSetInfo *) fcinfo->resultinfo;
InitMaterializedSRF(fcinfo, 0);
/* Make sure we get consistent view of the workers. */
LWLockAcquire(LogicalRepWorkerLock, LW_SHARED);
for (i = 0; i < max_logical_replication_workers; i++)
{ /* for each row */
Datum values[PG_STAT_GET_SUBSCRIPTION_COLS] = {0}; bool nulls[PG_STAT_GET_SUBSCRIPTION_COLS] = {0}; int worker_pid;
LogicalRepWorker worker;
memcpy(&worker, &LogicalRepCtx->workers[i], sizeof(LogicalRepWorker)); if (!worker.proc || !IsBackendPid(worker.proc->pid)) continue;
if (OidIsValid(subid) && worker.subid != subid) continue;
/* *Ifonlyasinglesubscriptionwasrequested,andwefoundit, *break.
*/ if (OidIsValid(subid)) break;
}
LWLockRelease(LogicalRepWorkerLock);
return (Datum) 0;
}
Messung V0.5 in Prozent
¤ Die Informationen auf dieser Webseite wurden
nach bestem Wissen sorgfältig zusammengestellt. Es wird jedoch weder Vollständigkeit, noch Richtigkeit,
noch Qualität der bereit gestellten Informationen zugesichert.0.30Bemerkung:
(vorverarbeitet am 2026-09-28)
¤
Die Informationen auf dieser Webseite wurden
nach bestem Wissen sorgfältig zusammengestellt. Es wird jedoch weder Vollständigkeit, noch Richtigkeit,
noch Qualität der bereit gestellten Informationen zugesichert.
Bemerkung:
Die farbliche Syntaxdarstellung und die Messung sind noch experimentell.