ingestion-pipeline-doctor-nodejs

bởi posthog

Tài liệu tham khảo nhanh về khung pipeline thu thập dữ liệu của PostHog và các tác nhân kiểm tra quy ước của nó.

npx skills add https://github.com/posthog/posthog --skill ingestion-pipeline-doctor-nodejs

Pipeline Doctor

Quick reference for PostHog's ingestion pipeline framework and its convention-checking agents.

Architecture overview

The ingestion pipeline processes events through a typed, composable step chain:

Kafka message
  → messageAware()
    → parse headers/body
    → sequentially() for preprocessing
    → filterMap() to enrich context (e.g., team lookup)
    → teamAware()
      → concurrentlyPerGroup(token:distinctId) for per-entity processing
      → gather()
      → pipeChunk() for chunk operations
      → handleIngestionWarnings()
    → handleResults()
  → handleSideEffects()
  → build()

See nodejs/src/ingestion/pipelines/analytics/joined-ingestion-pipeline.ts for the real implementation.

Key file locations

WhatWhere
Step typenodejs/src/ingestion/framework/steps.ts
Result typesnodejs/src/ingestion/framework/results.ts
Doc-test chaptersnodejs/src/ingestion/framework/docs/*.test.ts
Joined pipelinenodejs/src/ingestion/pipelines/analytics/joined-ingestion-pipeline.ts
Doctor agents.claude/agents/ingestion/
Test helpersnodejs/src/ingestion/framework/docs/helpers.ts

Which agent to use

ConcernAgentWhen to use
Step structurepipeline-step-doctorFactory pattern, type extension, config injection, naming
Result handlingpipeline-result-doctorok/dlq/drop/redirect, side effects, ingestion warnings
Compositionpipeline-composition-doctorBuilder chain, concurrency, grouping, branching, retries
Testingpipeline-testing-doctorTest helpers, assertions, fake timers, doc-test style

Quick convention reference

Steps: Factory function returning a named inner function. Generic <T extends Input> for type extension. No any. Config via closure.

Results: Use ok(), dlq(), drop(), redirect() constructors. Side effects as promises in ok(value, [effects]). Warnings as third parameter.

Composition: messageAware wraps the pipeline. handleResults inside messageAware. handleSideEffects after. concurrentlyPerGroup for per-entity work. gather before chunk steps.

Batching lifecycle hooks (BatchingPipeline beforeBatch/afterBatch): enrich-only. Hooks may enrich elements and batch context but must return exactly the elements they received — a count change is a broken invariant and feed() throws. Filtering belongs in sub-pipeline steps that return drop(). An empty feed() is a no-op (no hooks, no capacity). Details: nodejs/src/ingestion/framework/docs/14-batching.test.ts.

Fan-out/fan-in (fanOut(fn).via((sub) => …).fanIn(fn)): per-element sub-work with cardinality restored — one element fans out to N sub-elements (e.g. per-blob uploads), a regular sub-pipeline processes them (maxConcurrency on the sub concurrently block, retry on the per-sub step), and fan-in folds the OK results back into the parent. Reach for it over concurrently/concurrentlyPerGroup when the unit of concurrency is smaller than the element; hand-rolled p-limit/Promise.all inside a step is the tell. Sequencing is compile-time enforced (an unclosed stage cannot build). Sub-result contract: OK collected; DROP excludes the sub silently; DLQ fails the parent with aggregated reasons; REDIRECT is excluded with a warning — sub redirects never escape the stage. Sub-pipelines are context-agnostic: team/message data goes in the sub-element value, and context-gated surface (teamAware, handleIngestionWarnings, …) is uncallable. Fan-out/fan-in functions are cheap, synchronous, and named. Parents emit unordered as they complete. Details: nodejs/src/ingestion/framework/docs/17-fan-out-fan-in.test.ts.

Testing: Step tests call factory directly. Use consumeAll()/collectChunks() helpers. Fake timers for async. Type guards for result assertions. No any.

Running all doctors

Ask Claude to "run all pipeline doctors on my recent changes" to get a comprehensive review across all 4 concern areas.

Thêm skills từ posthog

error-tracking-hono
posthog
Theo dõi lỗi PostHog cho Hono
tuning-incremental-sync-config
posthog
Cấu hình của một đồng bộ nằm trên ExternalDataSchema và có thể được thay đổi bất kỳ lúc nào qua external-data-schemas-partial-update. Hầu hết các thay đổi đều không phá hủy (có hiệu lực vào lần đồng bộ tiếp theo), nhưng một số thay đổi (chuyển đổi sync_type, thay đổi khóa chính) yêu cầu xử lý cẩn thận để tránh làm hỏng dữ liệu đã đồng bộ.
playwright-test
posthog
Viết một bài kiểm tra playwright, đảm bảo nó chạy được và không bị lỗi không ổn định.
error-tracking-ruby
posthog
PostHog theo dõi lỗi cho Ruby
authoring-log-alerts
posthog
Tạo cảnh báo log hữu ích, ít nhiễu trên các dịch vụ trong một dự án PostHog. Sử dụng khi người dùng yêu cầu thiết lập cảnh báo cho log của họ, đề xuất các cảnh báo họ nên thêm,…
making-scenes-tab-aware
posthog
Guides converting PostHog frontend scenes to be tab aware for internal scene tabs. Use when adding or refactoring a `SceneExport` scene, fixing state leaking…
posthog-survey-creator
posthog
Tạo và cấu hình khảo sát trong PostHog thông qua hội thoại có hướng dẫn. Sử dụng kỹ năng này khi người dùng muốn tạo khảo sát, thu thập phản hồi người dùng, chạy…
authoring-scouts
posthog
Cách tạo, chỉnh sửa và điều chỉnh các scout PostHog Signals — các tác nhân theo lịch trình quét một dự án và viết báo cáo vào hộp thư đến Signals. Sử dụng khi người dùng…