ingestion-pipeline-doctor-nodejs

작성자: posthog

PostHog의 수집 파이프라인 프레임워크와 그 규칙 검사 에이전트에 대한 빠른 참조입니다.

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.

posthog의 다른 스킬

error-tracking-hono
posthog
PostHog 오류 추적 for Hono
tuning-incremental-sync-config
posthog
동기화의 구성은 ExternalDataSchema에 저장되며, external-data-schemas-partial-update를 통해 언제든지 변경할 수 있습니다. 대부분의 변경은 비파괴적이며(다음 동기화에 적용됨), 일부 변경(sync_type 전환, 기본 키 변경)은 동기화된 데이터 손상을 방지하기 위해 신중한 처리가 필요합니다.
playwright-test
posthog
플레이라이트 테스트를 작성하고, 실행이 잘 되며, 불안정하지 않도록 하세요.
error-tracking-ruby
posthog
PostHog Ruby 오류 추적
authoring-log-alerts
posthog
PostHog 프로젝트의 서비스에 유용하고 노이즈가 적은 로그 알림을 작성합니다. 사용자가 로그에 대한 알림 설정을 요청하거나 추가해야 할 알림을 제안할 때 사용하세요.
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
PostHog에서 안내 대화를 통해 설문조사를 생성하고 구성합니다. 사용자가 설문조사를 만들거나, 사용자 피드백을 수집하거나, 실행하려 할 때 이 스킬을 사용하세요.
authoring-scouts
posthog
PostHog Signals 스카우트를 작성, 편집 및 조정하는 방법 — 프로젝트를 스캔하고 Signals 인박스에 보고서를 작성하는 예약된 에이전트입니다. 사용자가…