How PostHog’s Data Warehouse Syncs External Sources: A Temporal-Driven Architecture

PostHog’s Data Warehouse uses Temporal schedules to orchestrate syncs between external sources (Stripe, Postgres, etc.) and ClickHouse, with declarative configuration managed via ExternalDataSource and ExternalDataSchema models.

The PostHog/posthog repository implements a robust data warehouse sync external sources mechanism that leverages Temporal for durable, fault-tolerant execution. This architecture cleanly separates schedule management from extraction logic, enabling reliable, resumable pipelines for ingesting third-party data into ClickHouse-backed tables.

Core Data Models That Drive Syncs

The synchronization system centers on two domain models defined in the data warehouse backend.

ExternalDataSource

The ExternalDataSource class represents the connection to an external system (e.g., Stripe, Postgres). Defined in products/data_warehouse/backend/models/external_data_source.py (lines 24-62), this model stores connection credentials, source type, and hierarchical links to its associated schemas.

ExternalDataSchema

The ExternalDataSchema class describes a single table or view that should materialize inside the Data Warehouse. Located in products/data_warehouse/backend/models/external_data_schema.py (lines 31-71), it stores the sync type (FULL_REFRESH, INCREMENTAL, CDC, or webhook) and schedule configuration. Each schema record maps to exactly one Temporal schedule.

How Temporal Scheduling Works

When a schema is created or updated, PostHog builds a Temporal Schedule that determines when to run the external-data-job workflow.

Building the Sync Schedule

The get_sync_schedule() function in products/data_warehouse/backend/data_load/service.py (lines 55-97) computes the execution time based on either a fixed sync_time_of_day or a jittered time derived from the sync_frequency_interval. This prevents the thundering herd problem by distributing syncs evenly across the day rather than overwhelming resources at midnight.

Converting to Temporal Objects

The to_temporal_schedule() method (lines 100-145 in products/data_warehouse/backend/data_load/service.py) instantiates a Temporal Schedule object containing a ScheduleActionStartWorkflow. This action launches the "external-data-job" workflow, passing an ExternalDataWorkflowInputs payload containing the team, schema, and source IDs.

Persisting and Triggering Schedules

The sync_external_data_job_workflow() function (lines 47-59) handles persistence logic:

  • If the schema is new (create=True), it calls create_schedule() and optionally triggers an immediate first run.
  • If the schema exists, update_schedule() rewrites the Temporal specification without interrupting ongoing executions.

To manually trigger a sync (e.g., from the UI), trigger_external_data_workflow() (lines 82-85) invokes the schedule’s trigger_schedule() method, causing Temporal to start the workflow on the next available tick.

CDC (Change Data Capture) Support

For CDC-enabled schemas, PostHog creates a separate source-level schedule that runs continuously at the tightest interval required by any CDC schema on that source.

Source-Level CDC Schedules

The get_cdc_extraction_schedule() function in products/data_warehouse/backend/data_load/service.py (lines 81-113) assembles a Temporal Schedule that invokes the "cdc-extraction" workflow. This schedule runs at the minimum sync_frequency_interval across all CDC schemas belonging to the source.

The sync_cdc_extraction_schedule() method (lines 123-163) computes the smallest interval across the source’s CDC schemas, creates or updates the schedule, or deletes it when no CDC schemas remain. This consolidates multiple CDC tables into one efficient extraction loop.

End-to-End Sync Flow

The complete data warehouse sync external sources process follows this execution path:

  1. User adds an external source → An ExternalDataSource record is saved.
  2. User defines schemas → ExternalDataSchema records are saved with a chosen sync_type.
  3. Model signals trigger scheduling → On each save, Django signals call sync_external_data_job_workflow(schema, create=True), which generates a Temporal Schedule via to_temporal_schedule() and persists it via create_schedule().
  4. Temporal execution → The Temporal engine triggers the "external-data-job" workflow at the scheduled time. The workflow runs the appropriate extractor and writes data into ClickHouse tables backing the DW.
  5. CDC streaming → For CDC schemas, the source-level CDC schedule runs the "cdc-extraction" workflow, continuously streaming changes into the DW.

All schedule state (paused status, memo, overlap policy) lives in Temporal; PostHog only records the schedule ID on the schema. The UI can pause or resume a schema by calling pause_external_data_schedule(schema.id) or unpause_external_data_schedule(schema.id).

Code Examples

Here are the practical implementations used to orchestrate syncs:


# Build or update a schedule for a schema (called from the API)

from products.data_warehouse.backend.data_load.service import sync_external_data_job_workflow

def create_or_update_schema(schema):
    # `create=True` on first insertion ensures the schedule is created

    sync_external_data_job_workflow(schema, create=True, should_sync=True)

# Manually trigger a sync (e.g., from a button in the UI)

from products.data_warehouse.backend.data_load.service import trigger_external_data_workflow

def run_now(schema):
    # Fires the Temporal schedule immediately

    trigger_external_data_workflow(schema)

# Enable CDC for a source – the system will auto-create a source-level schedule

from products.data_warehouse.backend.data_load.service import sync_cdc_extraction_schedule

def maybe_enable_cdc(source):
    # Called when a CDC schema is added/updated

    sync_cdc_extraction_schedule(source, create=False)

The actual extraction logic resides in posthog/temporal/data_imports/workflows/external_data_job.py, which implements the imperative data pulling logic separate from the declarative schedule management shown above.

Summary

  • PostHog’s Data Warehouse uses Temporal schedules to manage external source synchronization, separating schedule orchestration from extraction logic.
  • Two core models drive the system: ExternalDataSource (connection) and ExternalDataSchema (table configuration with sync type).
  • The sync_external_data_job_workflow() function in products/data_warehouse/backend/data_load/service.py bridges model changes to Temporal schedules via get_sync_schedule() and to_temporal_schedule().
  • CDC schemas trigger a separate source-level schedule through sync_cdc_extraction_schedule() to handle continuous change streaming at the minimum required interval.
  • Schedules can be paused, resumed, or manually triggered via trigger_external_data_workflow() without modifying underlying schema configuration.

Frequently Asked Questions

What is the difference between ExternalDataSource and ExternalDataSchema?

ExternalDataSource represents the connection to an external system (e.g., Stripe, Postgres) and is defined in products/data_warehouse/backend/models/external_data_source.py. ExternalDataSchema represents a specific table or view to materialize within the Data Warehouse, defined in products/data_warehouse/backend/models/external_data_schema.py, and includes the sync type and schedule configuration.

How does PostHog handle incremental vs full-refresh syncs?

The ExternalDataSchema model stores the sync_type (e.g., FULL_REFRESH, INCREMENTAL, CDC). The Temporal workflow receives this configuration via ExternalDataWorkflowInputs and executes the appropriate extraction strategy. Incremental syncs use timestamp or cursor-based filtering, while full-refresh syncs replace the entire table contents.

Can I manually trigger a data warehouse sync?

Yes. The trigger_external_data_workflow() function in products/data_warehouse/backend/data_load/service.py (lines 82-85) allows manual execution by calling the schedule’s trigger_schedule() method. This is typically invoked from the UI via the signals layer in products/signals/backend/views.py.

Where does the actual data extraction logic live?

While schedule management lives in products/data_warehouse/backend/data_load/service.py, the actual extraction workflows are implemented in posthog/temporal/data_imports/workflows/external_data_job.py. This separation ensures that the orchestration layer (Temporal schedules) remains decoupled from the imperative business logic of querying external APIs or databases.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →