Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 15 additions & 27 deletions docs/internal/cdc-buffered-ingress-runbook.md
Original file line number Diff line number Diff line change
Expand Up @@ -134,10 +134,9 @@ columns reports no position on that run, since none of its files carries the sta
residual buffer is merged once more and billed once; the next write lands with the statistic and the
position reads normally from then on.

**A run stands down while any delivery for the schema is still in the queue**: a previous attempt
of this same job, or a batch the retired legacy lane queued before its source was converted. Either
would write alongside whatever this run reads, and on the append lane that is a second copy of the
same history. Two scheduled runs cannot overlap on their
**A run stands down while any delivery for the schema is still in the queue**: a batch that a
previous attempt of this same job left there. It would write alongside whatever this run reads, and
on the append lane that is a second copy of the same history. Two scheduled runs cannot overlap on their
own: the v3 pipeline lock is held from the start of the workflow until the loader completes the
job. The window is a retried activity, which runs under the lock its own workflow already holds,
and a lock takeover, which hands the lock to a new job while the old one's batches are still
Expand All @@ -147,29 +146,18 @@ consumed. The next scheduled run picks the buffer up once the queue has drained.

## Leftover legacy state

Capture used to deliver some tables' changes itself, on what is now the retired legacy lane.
That lane paused each such table's schedule while it streamed, held a snapshotting table's changes as deferred runs in `sync_type_config["cdc_deferred_runs"]`, marked its source `cdc_ingest_mode: legacy`, and created job rows of its own.
Capture converts that state before every read (`cdc/legacy_conversion.py`), so nothing needs doing by hand:

- **A source still marked legacy** has every CDC table's buffer emptied, since it holds only copies of changes the legacy lane already delivered.
Each table is stamped `cdc_legacy_converted_at`, because it is current as of the conversion: without the stamp, its first sync would find no run that ever listed the buffer, take the buffer for expired, and re-snapshot the table (see "Buffer expiry — no partial recovery").
Each syncing table's schedule is rebuilt unpaused, because it is now the table's consumer.
A table set slower than its source's fastest table is sped up to it (`cdc_legacy_table_frequency_raised`), because that is how often the legacy lane delivered its changes.
The source is then marked `cdc_ingest_mode: buffered`, last, so a failure repeats the whole conversion on the next run.
Legacy batches still in the load queue land first, because the consumer stands down while any are in flight.
Until then, a scheduled run of one of its tables, such as one a sync frequency change unpaused, no-ops the tick (`cdc_buffered_waiting_for_legacy_conversion`), because reading those copies would load them a second time.
- **A table with deferred runs** snapshots again in the buffer, through the same pending reset a resync hands to capture (`cdc_legacy_deferred_runs_handed_to_reset`).
Nothing merges deferred runs anymore, so the reset pauses the table's schedule and cancels its running sync, and waits while either can still hand over (`cdc_reset_waits_for_running_sync`).
Then it resets the table, drops the deferred runs, empties its buffer, and starts the new snapshot.
Capture gives the table no buffered snapshot of its own before that, because the old sync could hand over into it without the deferred changes.
- **A job row a legacy capture run left Running** is failed once it is 30 minutes old and has no batches in the queue.

A rebuilt schedule is skipped where the pause is deliberate: the schema's status is `Paused`, an admin-triggered run holds it, the schema is halted and waits for Repair CDC, a pending reset holds it, or it has no sync frequency.
Every step logs (`cdc_legacy_source_converted`, `cdc_legacy_deferred_runs_handed_to_reset`, `cdc_stranded_capture_jobs_closed`); once none of them appears across the fleet, the module can go.

A source whose slot is gone does not capture at all, so it converts only after Repair CDC, which already resets every table and marks the source buffered.

`cdc_ingest_mode` must survive every API write. A PATCH that dropped it would make capture read the source as legacy and empty its unconsumed buffer.
Capture used to deliver some tables' changes itself, on the retired legacy lane.
The code that converted what that lane left behind is gone.
Every source that still read as legacy when it was removed is marked `cdc_broken`, so Repair CDC is its only way back.
Repair CDC resets every table, recreates the slot and resumes the table schedules, so no legacy state survives it.

Two leftovers remain in the data, and nothing reads them:

- `cdc_ingest_mode` in a source's `job_inputs`.
CDC setup and Repair CDC still write `buffered`, and the API still keeps the key on a PATCH.
A rollback to the release that converted legacy sources would read a source without it as legacy and empty its unconsumed buffer.
- `cdc_deferred_runs` in a table's `sync_type_config`.
A resync, a table-mode change and Repair CDC remove it.

## When a schedule stops firing

Expand Down
8 changes: 3 additions & 5 deletions products/warehouse_sources/backend/ad_hoc_sync.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,6 @@
from products.warehouse_sources.backend.temporal.data_imports.cdc.snapshot_lane import (
BUFFER_LANE,
cancel_sync_that_could_hand_over,
resnapshot_stays_in_buffer,
)


Expand Down Expand Up @@ -80,7 +79,7 @@
return bool(desc.schedule.state.paused)


def trigger_ad_hoc_sync(

Check warning on line 82 in products/warehouse_sources/backend/ad_hoc_sync.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

lint:complexity

`trigger_ad_hoc_sync` has cyclomatic complexity 13 (warn >10)

Check warning on line 82 in products/warehouse_sources/backend/ad_hoc_sync.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

`trigger_ad_hoc_sync` has cyclomatic complexity 13 (warn >10)
client: Client,
schema: ExternalDataSchema,
*,
Expand Down Expand Up @@ -126,10 +125,9 @@
# streaming, so ongoing CDC stays billable. The save must precede the workflow start so the
# source reloads cdc_mode="snapshot" instead of racing on stale "streaming".
if schema.is_cdc and schema.cdc_mode == "streaming":
# Decided while the table still streams. Without the marker, the next capture run would
# empty the buffer under the new snapshot, deleting changes an in-flight run wrote.
if resnapshot_stays_in_buffer(schema):
updates[CDC_SNAPSHOT_LANE_KEY] = BUFFER_LANE
# Without the marker, the next capture run would empty the buffer under the new snapshot,
# deleting changes an in-flight run wrote.
updates[CDC_SNAPSHOT_LANE_KEY] = BUFFER_LANE
updates["cdc_mode"] = "snapshot"
removes += ["cdc_last_log_position", "cdc_deferred_runs"]
extra_model_fields["initial_sync_complete"] = False
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -303,7 +303,7 @@
sensitive: set[str]


def get_nonsensitive_and_sensitive_field_names(fields: list[FieldType]) -> FieldSensitivitySplit:

Check warning on line 306 in products/warehouse_sources/backend/presentation/views/external_data_source/helpers.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

lint:complexity

`get_nonsensitive_and_sensitive_field_names` has cyclomatic complexity 11 (warn >10)

Check warning on line 306 in products/warehouse_sources/backend/presentation/views/external_data_source/helpers.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

`get_nonsensitive_and_sensitive_field_names` has cyclomatic complexity 11 (warn >10)
"""Classify source config field names as nonsensitive or sensitive.

Returns the field-name sets flattened across all nesting levels.
Expand Down Expand Up @@ -360,8 +360,8 @@
"cdc_lag_warning_threshold_mb",
"cdc_lag_critical_threshold_mb",
"cdc_consistent_point",
# Set by CDC setup, Repair CDC and capture, never by the API. Losing it on an unrelated PATCH
# would make capture convert the source again, which empties its unconsumed buffer.
# Set by CDC setup and Repair CDC, never by the API. A worker on a release that reads it treats a
# source without it as legacy and empties its buffer, so a PATCH must keep it for a rollback.
"cdc_ingest_mode",
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@
from products.warehouse_sources.backend.temporal.data_imports.cdc.broken import (
AUTO_DROPPED_LAG_REASON,
SELF_MANAGED_LAG_REASON,
broken_for_another_reason,
clear_recovered_self_managed_lag,
clear_slot_loss_markers,
mark_cdc_broken,
Expand All @@ -69,7 +70,6 @@
CDCSlotNotConfiguredError,
classify_cdc_error,
)
from products.warehouse_sources.backend.temporal.data_imports.cdc.legacy_conversion import convert_legacy_cdc_state
from products.warehouse_sources.backend.temporal.data_imports.cdc.load_resolution import has_engine_seq
from products.warehouse_sources.backend.temporal.data_imports.cdc.naming import (
CDC_EXTRACTION_WORKFLOW_ID_PREFIX,
Expand Down Expand Up @@ -510,25 +510,13 @@
return True

def _prepare_buffer(self) -> None:
"""Convert leftover legacy state and start pending snapshots in the buffer, before the WAL read."""
assert self.source is not None and self.adapter is not None
convert_legacy_cdc_state(
self.source,
self.cdc_schemas,
ingest_mode=self.adapter.parse_cdc_config(self.source).ingest_mode,
logger=self.log,
)

"""Start pending snapshots in the buffer, before the WAL read."""
for schema in self.cdc_schemas:
if not captures_to_buffer(schema):
# No lane writes this table mode, so the buffer could never deliver its changes.
self._schema_log(schema).warning("cdc_table_mode_not_captured", cdc_table_mode=schema.cdc_table_mode)
continue
# A table with deferred runs gets no buffered snapshot here: its old sync could still hand over
# into that buffer without its deferred changes. The reset the conversion staged restarts the
# snapshot instead, and holds the table out of capture until the old sync stops. Once that
# reset has run, which can be before this read, the table's changes belong in the buffer.
if snapshot_can_start_in_buffer(schema) and not (schema.sync_type_config or {}).get("cdc_deferred_runs"):
if snapshot_can_start_in_buffer(schema):
self._start_snapshot_in_buffer(schema)
self._buffered_table_names.add(schema.name)
self.log.info("cdc_buffered_ingress_active", buffered=sorted(self._buffered_table_names))
Expand Down Expand Up @@ -756,7 +744,7 @@

return heartbeat

def _read_wal_loop(self) -> None:

Check warning on line 747 in products/warehouse_sources/backend/temporal/data_imports/cdc/activities.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

lint:complexity

`_read_wal_loop` has cyclomatic complexity 12 (warn >10)

Check warning on line 747 in products/warehouse_sources/backend/temporal/data_imports/cdc/activities.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

`_read_wal_loop` has cyclomatic complexity 12 (warn >10)
"""Read WAL events with periodic micro-batch flushes, bounded per peek.

Each pass peeks at most ``CDC_MAX_CHANGES_PER_READ`` changes so a large backlog can't
Expand Down Expand Up @@ -1529,7 +1517,7 @@


@activity.defn
def cleanup_orphan_slots_activity() -> None:

Check warning on line 1520 in products/warehouse_sources/backend/temporal/data_imports/cdc/activities.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

lint:complexity

`cleanup_orphan_slots_activity` has cyclomatic complexity 29 (warn >10)

Check warning on line 1520 in products/warehouse_sources/backend/temporal/data_imports/cdc/activities.py

View workflow job for this annotation

GitHub Actions / Python code quality (depot-ubuntu-24.04)

`cleanup_orphan_slots_activity` has cyclomatic complexity 29 (warn >10)
"""Safety-net sweeper: clean up orphaned CDC slots and monitor WAL lag.

1. For deleted/inactive PostHog-managed sources → drop slot + publication
Expand Down Expand Up @@ -1705,17 +1693,20 @@
elif cdc_config.management_mode == "self_managed":
# Customer owns the slot: surface the broken state but keep the schedule running
# and never drop — the lag may recover once they reduce load on the source.
# A marker with another reason stays. The lag marker allows Resume CDC and clears
# itself once the lag drops, which would lift a stop that only Repair CDC may lift.
try:
mark_cdc_broken(
source,
SELF_MANAGED_LAG_REASON,
f"Change data capture replication lag exceeded {critical_threshold_mb} MB. "
f"This slot is self-managed, so PostHog did not drop it — reduce load or WAL "
f"retention on the source database, or it may invalidate the slot and "
f"require a full re-sync.",
pause=False,
lag_mb=round(lag_mb, 1),
)
if not broken_for_another_reason(source, SELF_MANAGED_LAG_REASON):
mark_cdc_broken(
source,
SELF_MANAGED_LAG_REASON,
f"Change data capture replication lag exceeded {critical_threshold_mb} MB. "
f"This slot is self-managed, so PostHog did not drop it — reduce load or WAL "
f"retention on the source database, or it may invalidate the slot and "
f"require a full re-sync.",
pause=False,
lag_mb=round(lag_mb, 1),
)
except Exception:
source_log.exception("failed_to_mark_self_managed_broken")
metrics.get_sweeper_source_errors_metric().add(1)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,25 @@ def mark_cdc_broken(
log.warning("cdc_marked_broken", schemas=len(cdc_schemas), newly_broken=len(newly_broken), paused=pause)


def broken_for_another_reason(source: ExternalDataSource, reason: str) -> bool:
"""Whether a table that `mark_cdc_broken` would mark already holds a marker with a different reason.

A table whose sync is off keeps the marker it had, so it must not count.
"""
return (
Comment thread
coderabbitai[bot] marked this conversation as resolved.
ExternalDataSchema.objects.filter(
team_id=source.team_id,
source=source,
sync_type=ExternalDataSchema.SyncType.CDC,
should_sync=True,
sync_type_config__has_key="cdc_broken",
)
.exclude(deleted=True)
.exclude(sync_type_config__cdc_broken__reason=reason)
.exists()
)


def clear_recovered_self_managed_lag(source: ExternalDataSource) -> int:
"""Lift the ``critical_lag_self_managed`` marker once the slot's lag is back under the warning threshold.

Expand Down
Loading
Loading