How PostHog's Event Capture Endpoint Processes and Ingests Events
The public /api/event/ endpoint validates incoming requests and forwards a compact JSON payload to the Rust-based capture-rs service, which writes to Kafka; asynchronous workers then resolve person profiles, apply event pipelines, and persist the fully processed data to ClickHouse.
When you track events using the PostHog JavaScript SDK or the raw HTTP API, your data enters a high-throughput ingestion pipeline designed to handle billions of events daily. The PostHog event capture endpoint—exposed publicly as /api/event/ (and internally forwarded to the capture service)—is a thin Django wrapper that delegates heavy processing to specialized downstream services. This article traces the exact code path from the initial HTTP request through to final storage, referencing specific files and functions from the PostHog/posthog repository.
Django Entry Point: EventViewSet and Validation
The ingestion journey begins in posthog/api/event.py inside the EventViewSet class. When a client POSTs to /api/event/, the view parses the request, verifies the API key against the team, and validates required fields.
According to the source code, the view utilizes validation helpers in posthog/api/utils.py (specifically capture_endpoint_invalid_payload) to ensure the payload contains mandatory keys: event, distinct_id, and properties. Malformed payloads immediately raise a ValidationError before any downstream processing occurs. Once validated, the view determines whether the request represents a standard analytics event or a session recording snapshot, setting internal flags accordingly.
Internal Capture Forwarding: The capture_internal Bridge
After validation, the Django view delegates to capture_internal defined in posthog/api/capture.py. This function acts as the bridge between the Python application layer and the Rust-based ingestion service.
The capture_internal function performs three critical operations:
- Payload Normalization: It calls
prepare_capture_internal_payloadto inject internal metadata flags (capture_internal=True,$process_person_profile), normalize timestamps to UTC ISO format (datetime.now(UTC).isoformat()), and guarantee the presence of asent_atfield. - Smart Routing: Based on the
event_name, the function routes session recording events (identified by$snapshotor$performance_event) to a dedicated replay endpoint, while standard analytics events flow toNEW_ANALYTICS_CAPTURE_ENDPOINT. - Reliable HTTP POST: The function uses
internal_requests_session()(configured with an HTTPAdapter and retry logic) to POST the JSON payload toCAPTURE_INTERNAL_URL. The request retries up to three times on 5xx errors before failing.
Crucially, the Django endpoint sets process_person_profile=False during this handoff—person profile resolution is intentionally deferred to downstream workers to minimize latency on the hot path.
Kafka Ingestion via Capture-RS
Once capture_internal successfully POSTs the payload, the capture-rs service (a Rust-based component referenced in posthog/settings/ingestion.py) takes over. This service writes the event directly to the events Kafka topic, preserving ordering guarantees per team ID.
This architectural separation allows the Django application to remain stateless and horizontally scalable while the Rust service handles the high-throughput, low-latency requirement of Kafka producer connections.
Async Processing and ClickHouse Persistence
Events sitting in the events Kafka topic are consumed by async workers defined in posthog/tasks/event_ingestion.py (Celery) and workflows in posthog/temporal/. These workers execute the heavy lifting of the ingestion pipeline:
- Person Resolution: The worker calls
get_persons_by_distinct_idsfromposthog/models/person/util.pyto map thedistinct_idto an existing Person record or create a new one, applying the profile updates that were deferred earlier. - Event Pipelines: The system runs enrichment and filtering logic defined in modules like
posthog/event_usage.py, applying rate limiting fromposthog/rate_limit.pywhere necessary. - ClickHouse Insertion: Finally, the worker executes the SQL defined in
posthog/models/event/sql.py(specificallyINSERT_EVENT_SQL), persisting the event with columns forteam_id,eventname,propertiesJSON,timestamp, and other metadata into the ClickHouseeventstable.
Subsequent reads—whether through the /api/event/ list endpoint or HogQL queries—hit this ClickHouse table directly, leveraging the indexing and query optimization defined in the same SQL modules.
Observability and Monitoring
The entire pipeline is instrumented with Prometheus counters and OpenTelemetry tracing. Key metrics include capture_internal_event_submitted_counter (tracking successful forwards to capture-rs) and EVENT_VALUES_COUNTER (monitoring event volume), defined in posthog/api/metrics.py. Spans are created using tracer.start_as_current_span to provide distributed tracing across the Django, Kafka, and worker boundaries.
Practical Example: Sending an Event
Below is the minimal HTTP request structure that the capture_internal function forwards to the Rust service:
POST https://us.posthog.com/api/event/
Content-Type: application/json
User-Agent: posthog-python/1.23.0
{
"api_key": "phc_1234567890abcdef",
"event": "signup",
"properties": {
"plan": "enterprise",
"referrer": "google"
},
"distinct_id": "user-42",
"timestamp": "2024-04-25T12:34:56.789Z"
}
Using the official Python SDK:
from posthog import Posthog
client = Posthog(
project_api_key="phc_1234567890abcdef",
host="https://us.posthog.com"
)
client.capture(
distinct_id="user-42",
event="signup",
properties={"plan": "enterprise", "referrer": "google"},
)
Both methods trigger the same chain: validation in EventViewSet, forwarding via capture_internal, Kafka ingestion via capture-rs, and eventual persistence to ClickHouse.
Summary
- The Django EventViewSet (
posthog/api/event.py) validates API keys and payloads, rejecting malformed requests immediately. capture_internal(posthog/api/capture.py) normalizes payloads and forwards them via HTTP to the Rust-based capture service, with automatic retries and separate routing for session recordings.- Capture-rs writes events to the
eventsKafka topic, decoupling the web tier from the data layer. - Async workers (Celery/Temporal) consume from Kafka, resolve person profiles using
posthog/models/person/util.py, and insert into ClickHouse viaposthog/models/event/sql.py. - Full observability is provided through Prometheus counters and OpenTelemetry spans throughout the pipeline.
Frequently Asked Questions
What is the difference between the /capture and /api/event/ endpoints?
In the PostHog codebase, /api/event/ is the public REST endpoint exposed to SDKs, while internal logic treats this as the capture endpoint. The Django view in posthog/api/event.py handles these POST requests and immediately delegates to the capture_internal function in posthog/api/capture.py, effectively making them the same entry point in modern versions of the platform.
How does PostHog handle person profile resolution during ingestion?
Person profile resolution is intentionally asynchronous. The Django capture endpoint sets process_person_profile=False when calling capture_internal, ensuring the HTTP response returns quickly. Later, workers in posthog/tasks/event_ingestion.py read from Kafka and call get_persons_by_distinct_ids from posthog/models/person/util.py to link events to Person records or create new profiles without blocking the initial request.
What happens if the capture-rs service is temporarily unavailable?
The capture_internal function implements retry logic using an HTTPAdapter that retries up to three times on 5xx errors. If the service remains unavailable after retries, the Django view returns an error to the client, allowing SDKs to queue and retry events client-side according to their own retry policies.
Why does PostHog use a Rust service for Kafka ingestion instead of writing directly from Django?
The Rust-based capture-rs service isolates the high-throughput, low-latency requirements of Kafka production from the Django application server. This separation allows the Python layer to remain stateless and horizontally scalable while the optimized Rust binary handles connection pooling, batching, and guaranteed delivery to Kafka, preventing ingestion bottlenecks during traffic spikes.
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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →