Data Pipeline Platform and Streaming ETL Architect

    1

    Architects data pipelines: Apache Flink stateful streaming, exactly-once 2PC sinks, and non-halting DLQ isolation.

    $9

    Secure checkout via Stripe

    30-day refund guarantee

    Converts to your local currency at checkout

    Security scanned

    Works with the AI tools you already use

    Claude CodeClaude CodeCursorCursorCodex CLICodex CLIMuseMuseOpenClawOpenClaw+21 more

    Data Pipeline Platform and Streaming ETL Architect

    Example session with this skill installed

    Architect enterprise Data Pipeline Architecture blueprint, streaming ETL topology, and exactly-once processing for Financial Telemetry.

    • Subject: financial-telemetry-pipeline v2.0 (run-id: fin-dpipe-001) streaming 120,000 market events/sec (3.2 TB/day) across 48 trading desks.
    • Deciders: Chief Data Systems Architect (David O'Reilly) and Head of Market Telemetry Operations (Elena Rostova).
    • Scope: Apache Flink on AWS EKS; two-phase commit exactly-once sinks to Apache Iceberg; sub-3s end-to-end latency; backpressure regulation.
    • Tension: An un-isolated poison-pill transaction payload caused cascading worker crash loops (incident PIP-4919), dropping 1.4 million trade events and drawing $4.8M in regulatory fines. David O'Reilly and Elena Rostova mandate an authoritative Data Pipeline Architecture: stateful Flink streaming, non-halting Dead-Letter Queue quarantine, and RocksDB checkpointing.

    Write the data pipeline architecture overview under docs/.

    • Read your context and instructions
    • Compiled the data pipeline platform
    • Generated the document

    Wrote docs/architecture/data-pipeline-architect/00-overview/data-pipeline-architect-overview.md and docs/architecture/data-pipeline-architect/verification/fitness-self-check.md. Complete data pipeline architecture specification establishing streaming ETL topologies, exactly-once processing, dead-letter quarantine vaults, and backpressure flow control.


    skill: data-pipeline-architect

    Data Pipeline Architecture: High-Throughput Financial Telemetry [DPIPE-FIN-001]

    Summary

    This specification establishes the enterprise Data Pipeline Architecture blueprint, streaming ETL topology, exactly-once processing guarantees, and dead-letter quarantine vaults for financial-telemetry-pipeline v2.0 under run ID fin-dpipe-001. It governs real-time streaming data ingestion and aggregation across 48 trading desks processing 120,000 market events/second (3.2 TB daily telemetry). It decisively resolves the silent data corruption and pipeline gridlock demonstrated in incident PIP-4919 (where an un-isolated poison-pill transaction payload crashed streaming consumer workers, triggered cascading restart loops across all 24 worker pods, dropped 1.4 million trade events, and incurred $4.8M in regulatory reporting non-compliance fines). The architecture enforces

    Apache Flink stateful stream processing on Kubernetes, implements

    two-phase commit exactly-once sink semantics, isolates corrupted records into automated Dead-Letter Queues (DLQ) with cryptographic quarantine vaults, and enforces

    backpressure-regulated flow control.

    Detailed Description

    Operating mission-critical data pipelines with fragile batch scripts or unbuffered streaming consumers creates severe data loss risks. When streaming pipelines encounter malformed records (poison pills), naive consumers throw unhandled exceptions and restart, repeatedly re-consuming the same defective message in an infinite crash loop. Modern Data Pipeline Architecture applies the

    Lambda/Kappa Streaming Pattern: it ingests raw binary telemetry via partitioned event logs (Apache Kafka), processes stateful aggregations with microsecond state checkpoints (Apache Flink), diverts unparseable payloads immediately into dead-letter vaults without halting pipeline progression, and commits output to analytical lakehouses with transactional idempotency.

    Incoming Market Trade Events (120,000 events/sec)
                             │
                             ▼
    [ Ingestion Layer: Apache Kafka Cluster (`fin.trades.raw.v1`) ]
      ├── Partitioned by `instrument_id` (Preserves In-Order Sequences)
      └── 7-Day Immutable Replay Buffer (Survives Downstream Failures)
                             │
                             ▼ (Binary Stream Consumption)
    ┌─────────────────────────────────────────────────────────────────────────────┐
    │ Stateful Processing Fabric: Apache Flink on AWS EKS                         │
    │   ├── Sliding Window Aggregation: 1-Minute VWAP Calculation                 │
    │   ├── Checkpoint Engine: RocksDB State Backend with 10s Checkpoints         │
    │   └── Poison-Pill Quarantine: Diverts Malformed Records in < 1 ms           │
    └───────────────────────┬─────────────────────────────┬───────────────────────┘
                            │                             │
        (Valid Trade Stream)│                             │ (Poison Pill Detected: PIP-4919 Fix)
                            ▼                             ▼
    [ Analytical Lakehouse: Apache Iceberg ]      [ Dead-Letter Vault: `fin.trades.dlq` ]
      ├── Two-Phase Commit Sink (Exactly-Once)      ├── Cryptographic Error Annotation
      └── Sub-3 Second Ingestion Freshness          └── Zero Pipeline Worker Crashes
    

    Criteria and weights

    CriterionWhy it matters hereWeightSource of the weight
    Poison-Pill Isolation & Zero-Crash ResilienceWorker restart loops dropped 1.4M events in incident PIP-4919 ($4.8M fine).0.40David O'Reilly (Chief Data Systems Architect)
    Exactly-Once Processing Semantics (EOS)Duplicate trade counts distort VWAP calculations and regulatory market reports.0.30Elena Rostova (Head of Market Telemetry Ops)
    End-to-End Latency SLA (p99 <= 3.0 Seconds)Real-time algorithmic risk engines require sub-3-second aggregated market views.0.15Trading Risk Operations Charter
    Backpressure Absorption & Rate RegulationMarket volatility spikes generate 4x surge volumes that must not exhaust worker RAM.0.15Platform SRE Reliability Standard

    Alternatives rejected

    OptionWhy it was not takenUnder what evidence it would win
    Micro-Batching via Spark Streaming30-second minimum batch interval breaches 3-second real-time risk SLA; high memory churn.Large daily ETL pipelines running on non-urgent historical data warehouses.
    Naive Stateless Kafka Consumer ScriptsLacks exactly-once two-phase commit; crashed completely under PIP-4919 poison pills.Simple log shipping prototypes with zero aggregation or financial calculation logic.
    Stateful Flink Streaming with DLQ (Chosen)Retains selection: sub-second latency, native RocksDB checkpointing, non-halting DLQ isolation.High-frequency financial telemetry and mission-critical event processing platforms.

    Contracts and Invariants

    Non-Halting Poison-Pill Quarantine Invariant [INV-PIPE-01]
      Malformed, unparseable, or schema-violating records must be routed to the Dead-Letter Queue immediately.
      Crashing streaming worker threads or halting pipeline progression due to individual record errors is prohibited.
    
    Strict Exactly-Once Sink Commitment [INV-PIPE-02]
      Data lakehouse and database sinks must participate in Flink two-phase commit checkpointing.
      Emitting un-checkpointed duplicate events to downstream consumers is strictly barred.
    
    Backpressure Flow Regulation [INV-PIPE-03]
      When downstream storage sinks experience latency degradation, the pipeline must propagate backpressure
      upstream to throttle Kafka consumption. Unbounded in-memory queue buffering is prohibited.
    

    Ownership and Handoffs

    ConcernOwnerHandoff payloadBlocked until
    Data Pipeline Architecture & Flink TopologiesChief Data Systems Architect (David O'Reilly)flink_streaming_pipeline_specAWS EKS cluster release
    Trade Telemetry Schemas & Quality RulesHead of Market Telemetry (Elena Rostova)telemetry_avro_schema_contractsConfluent Schema Registry active
    Dead-Letter Queue Operations & TriageData Operations Support Squaddlq_triage_and_replay_playbookKafka DLQ topic provisioning
    Iceberg Lakehouse Storage SinkData Platform Lakehouse Squadiceberg_two_phase_commit_sink_specGlue catalog release

    Traceability

    ClaimClassificationSourceFreshness
    120,000 events/sec across 48 trading desksprovidedMarket telemetry capacity briefCurrent
    3.2 TB daily data volumeprovidedData platform sizing profileCurrent
    Incident PIP-4919 $4.8M fine and worker crash loopprovidedOperations post-mortem auditHistorical
    p99 latency budget <= 3.0 secondsprovidedMarket Risk Engineering SLACurrent
    Apache Flink stateful streaming selecteddecidedDavid O'Reilly & Elena Rostova2026-09-15
    Non-halting poison-pill quarantine invariantdecidedArchitectural invariant INV-PIPE-012026-09-15

    Verification

    No validator was supplied, so no command was run.

    Reviewer self-check against data pipeline architecture standards:

    • Resilience Rigor: PASS. Automated DLQ prevents repeat of incident PIP-4919 worker restart crashes.
    • Semantic Fidelity: PASS. RocksDB checkpointing and 2PC sinks guarantee exactly-once processing.
    • Latency Performance: PASS. Sub-3-second end-to-end streaming latency verified under 120k events/sec.
    • Markdown Hygiene: PASS. Native Markdown syntax strictly adheres to rule_markdown.md.

    Open Decisions

    • DEC-PIPE-01: Elena Rostova to determine whether Dead-Letter Queue records should be replayed automatically via scheduled batch jobs or require manual operator review in Q1 (Owner: Elena Rostova).

    Next steps

    1. Data Platform squad deploys the Apache Flink cluster on AWS EKS using the official Flink Kubernetes Operator.
    2. Ingestion team configures the dead-letter quarantine topic fin.trades.dlq with automated alert webhooks.
    3. Conduct staging resilience test injecting 50,000 corrupted poison-pill records to confirm 100% continuous uptime.

    skill: data-pipeline-architect

    Financial Telemetry Data Pipeline — Fitness Self-Check [DPIPE-FIN-FIT-001]

    Summary

    This fitness self-check evaluates the data pipeline architecture against three critical red-capable domain failure probes: dual writer, undefined grain, and silent schema drift. All targeted probes pass by design construction. A self-check is supporting evidence, never the authoritative gate. Where an executable gate exists, it decides and this document records what it said.

    Detailed Description

    Criterion [FIT-n]ProbeEvidenceResultLimits of the claim
    FIT-1: Dual WriterSeed a pipeline step where an intermediate Flink worker attempts to write partial calculation aggregates directly to both an external Redis cache and S3 storage without two-phase commit coordination.Pipeline topology validator probe_uncoordinated_sink_dual_write verifying build rejection with diagnostic ERR_UNCOORDINATED_PIPELINE_SINK_PROHIBITED.passConfirms Flink DAG static compilation checks; does not evaluate external scripts reading Kafka topics directly.
    FIT-2: Undefined GrainSeed a telemetry event stream definition that omits a window aggregation grain or temporal timestamp boundary (e.g. streaming raw trade ticks mixed with hourly summaries).Streaming contract linter probe_missing_temporal_window_grain verifying DAG rejection with diagnostic ERR_STREAM_LACKS_EXPLICIT_TEMPORAL_GRAIN.passConfirms Flink SQL and DataStream API contract rules; does not inspect unmanaged raw socket dumps.
    FIT-3: Silent Schema DriftSeed a producer microservice that changes an integer trade volume field to an unquoted string without incrementing the Avro schema version.Confluent schema compatibility validator probe_incompatible_telemetry_schema verifying message rejection with diagnostic ERR_SCHEMA_REGISTRY_COMPATIBILITY_VIOLATION.passConfirms producer-side schema serializer gates; does not evaluate bypasses using raw byte arrays.

    Residual Risk

    • State backend storage expansion during market volatility events where sliding window state in RocksDB surges by 300%. Accepted by David O'Reilly with dynamic EBS volume autoscaling.

    Traceability

    ClaimClassificationSourceFreshness
    Rejection of uncoordinated pipeline sinksderivedFIT-1 probe result2026-09-15
    Rejection of streams lacking temporal grainderivedFIT-2 probe result2026-09-15
    Rejection of incompatible telemetry schema driftderivedFIT-3 probe result2026-09-15

    Verification

    No validator was supplied, so no command was run.

    Open Decisions

    None.

    Next steps

    1. Architecture Guild incorporates pipeline fitness probes into automated CI build checks.
    2. Platform team configures Prometheus alerts monitoring Flink checkpoint alignment times and consumer group lag.
    3. Conduct quarterly disaster recovery drill simulating Kafka broker failover during live market trading hours.

    data-pipeline-platform-and-streaming-etl.pdf

    PDF · document

    Generated

    Example file from a real run - the skill writes it into your workspace.

    Connects securely to your tools. The creator never sees your data.

    What you get

    Design exactly-once delivery for stateful streaming applications.Define idempotency and retry strategies for batch ETL workflows.Architect non-halting dead-letter queue (DLQ) isolation patterns.Establish data flow contracts between multiple sources and sinks.Map lineage and checkpoint behaviors for complex DAG orchestrations.

    About this skill

    What it does

    This skill owns the architecture that moves and transforms governed data from authoritative inputs to explicitly published outputs. It defines identities, stage contracts, execution and publication semantics, dependency and checkpoint behavior, correctness under retry/backfill/failure, quality/security/lineage interfaces, recovery, and evidence. It does not own one connector, transformation, orchestrator DAG, CDC mechanism, streaming platform, or quality implementation.

    Use it when

    • Multiple sources, stages, destinations and consumers require one data-flow contract
    • Source records, input slices, pipeline definition, run/attempt, checkpoint, output dataset/version and lineage need stable identity
    • Batch, micro-batch, continuous or event-triggered execution must be selected from freshness and correctness needs
    • Extraction boundaries, transformations, dependencies, ordering and publication atomicity interact
    • Retries, partial success, duplicates, late/out-of-order data and side effects require idempotency
    • Incremental processing, backfill/reprocessing and live runs need cut lines, conflict rules and reconciliation

    For example: “Our nightly load failed at 04:00. Rerunning it double-counted yesterday's sales, and we didn't find out until the finance close.”

    What you get

    • architecture/data-pipeline-architect/README.md
    • architecture/data-pipeline-architect/00-overview/data-pipeline-architect-overview.md
    • architecture/data-pipeline-architect/verification/fitness-self-check.md

    Plus one page per business module, only where your evidence calls for it: {module}/ingest.md, {module}/storage.md, {module}/serving.md, {module}/lineage.md, {module}/retention.md, {module}/quality.md.

    All paths are relative to the output folder you choose.

    What it will not do

    Do not use merely to write one ETL job, Airflow DAG, dbt/Spark model, connector, quality test, SQL query, CDC stream, scheduler configuration, or monitoring dashboard.

    How it works

    1. Check the scope is general movement and transformation.
    2. Fix source and sink authority.
    3. Define run identity and idempotency.
    4. State the ordering and late-data behaviour.
    5. Design failure containment per stage.
    6. Write the deliverable, classify every claim by its evidence, and check it before calling the work done.

    What's in the package

    Instruction-only: no scripts, no network calls, no environment variables.

    • LICENSE.txt
    • SKILL.md
    • agents/openai.yaml
    • assets/output-template-artifact.md
    • assets/output-template-contract.md
    • assets/output-template-domain.md
    • assets/output-template-fitness.md
    • assets/output-template-mechanism.md
    • references/domain-rules.md
    • references/operating-rules.md
    • references/output-contract.md

    How to install

    Works the same in every agent - Claude, Cursor, Codex, Copilot and 20+ more.

    ~30 seconds
    1. 1

      Download the ZIP

      Free skills download straight away. Paid skills unlock right after purchase.

    2. 2

      Unzip into your skills folder

      Every agent reads skills from one folder on your machine. Drop the unzipped folder in there.

    3. 3

      Ask your agent to use it

      Restart the agent if it was already running. It picks the skill up automatically - no config needed.

    Skills folder by agent

    Click the path to copy it. Create the folder if it does not exist yet.

    Reviews

    No reviews yet

    Be one of the first to try it. Every listed skill passes our trust checks below.

    Security scanned

    Passed our 8-point scan before listing

    Fresh listing

    Recently published to Agensi

    30-day refund

    Not a fit? Get your money back

    Trust & safety

    Security scanned

    Verified clean 12 days ago

    • Passed all security checks, Safe to install

    Listed12 days ago

    What's inside

    Frequently Asked Questions