CommunityRecherche & Datenanalysegithub.com

Unknown-333/building-kafka-consumers

Build reliable Apache Kafka consumers and producers — consumer groups and partition assignment, offset commit strategy, at-least-once vs exactly-once, idempotent/transactional producers, rebalancing, and dead-letter handling. Use when writing Kafka consumers/producers, configuring offset commits or consumer groups, tuning throughput, or handling rebalances and poison messages.

Was ist building-kafka-consumers?

building-kafka-consumers is a Claude Code agent skill that build reliable Apache Kafka consumers and producers — consumer groups and partition assignment, offset commit strategy, at-least-once vs exactly-once, idempotent/transactional producers, rebalancing, and dead-letter handling. Use when writing Kafka consumers/producers, configuring offset commits or consumer groups, tuning throughput, or handling rebalances and poison messages.

Funktioniert mit~Claude Code~Codex CLI~Cursor
npx skills add https://github.com/Unknown-333/awesome-data-engineering-skills/tree/main/skills/building-kafka-consumers

Installed? Explore more Recherche & Datenanalyse skills: obra/superpowers, affaan-m/quarkus-verification, affaan-m/uspto-database · View all 6 →

In Ihrer bevorzugten KI fragen

Öffnet einen neuen Chat, in dem dieser Agent-Skill bereits geladen ist.

Dokumentation

Building Kafka Consumers

When to use

  • Writing or debugging Kafka consumers/producers.
  • Choosing offset-commit strategy and delivery guarantees.
  • Tuning consumer-group parallelism, rebalancing, or dead-letter handling.
  • Do NOT use for stream processing/windowing (use processing-streaming-data).

Workflow

- [ ] Size partitions to target parallelism (consumers <= partitions)
- [ ] Commit offsets AFTER successful processing
- [ ] Make the sink idempotent (upsert by event key)
- [ ] Handle rebalances (commit on revoke, avoid long poll gaps)
- [ ] Route poison messages to a dead-letter topic
  1. Partitions cap parallelism. A consumer group scales out only up to the partition count; extra consumers sit idle. Choose partitions for peak throughput.
  2. Commit after processing. Commit offsets once the work is durably done, not before — committing early loses messages on a crash.
  3. Idempotent sink. At-least-once means duplicates on retry; upsert by a stable event key so reprocessing is harmless.
  4. Rebalances happen. Commit on partition revoke and keep poll() intervals under max.poll.interval.ms so the broker doesn't evict the consumer.
  5. Poison messages go to a dead-letter topic with the error, so one bad record doesn't block the partition.

Patterns

Manual commit after processing:

consumer = KafkaConsumer("orders", group_id="etl",
                         enable_auto_commit=False,
                         max_poll_records=500)
for msg in consumer:
    try:
        upsert(process(msg))          # idempotent by key
        consumer.commit()             # commit only after success
    except PoisonError:
        send_to_dlq(msg)
        consumer.commit()             # skip the bad record

Idempotent / transactional producer — set enable.idempotence=true (dedupes retries) and use transactions for read-process-write exactly-once across topics.

Throughput tuning — increase max.poll.records, fetch.min.bytes, and process in batches; keep processing fast to avoid rebalance eviction.

Common pitfalls

  • Auto-commit + slow processing — offsets advance before work completes; a crash drops messages. Prefer manual commit after processing.
  • More consumers than partitions — the extras idle; repartition to scale.
  • Long processing between polls — exceeds max.poll.interval.ms and triggers endless rebalances; process in bounded batches or use a background worker.
  • No dead-letter path — one poison message blocks the whole partition.
  • Relying on exactly-once without an idempotent sink — any at-least-once hop reintroduces duplicates.

Individual skills in this repo

This repo contains 9 individual skills — each has its own dedicated page.

Unknown-333/authoring-airflow-dags

Write production-grade Apache Airflow DAGs using the TaskFlow API — idempotent tasks, correct scheduling and catchup, retries/SLAs, connections/variables, and avoiding top-level code. Use when creating or reviewing Airflow DAGs, scheduling pipelines, wiring task dependencies, configuring retries/backfills, or fixing non-idempotent tasks.

Unknown-333/building-dagster-assets

Build Dagster pipelines using software-defined assets — asset dependencies, partitions, resources and IO managers, asset checks, and schedules/sensors. Use when creating Dagster assets or jobs, modeling data as assets, adding partitions or backfills, wiring resources/IO managers, or migrating from task-based orchestration to assets.

Unknown-333/building-dbt-models

Build well-structured dbt models — staging/intermediate/marts layers, ref() and source(), materializations, and incremental models with the right strategy. Use when creating or refactoring dbt models, choosing table vs view vs incremental, structuring a dbt project, or writing incremental logic.

Unknown-333/building-feature-pipelines

Build ML feature pipelines and feature stores — point-in-time-correct joins to avoid label leakage, offline/online parity, feature freshness and backfills, and materialization with tools like Feast. Use when engineering features for ML, preventing train/serve skew or data leakage, building a feature store, or backfilling historical features for training.

Unknown-333/building-iceberg-tables

Design and operate Apache Iceberg tables — partitioning and hidden partitioning, partition/schema evolution, snapshots and time travel, compaction and small-file cleanup, and MERGE/upsert for lakehouse tables on Spark, Flink, Trino, or Snowflake. Use when creating or maintaining Iceberg tables, choosing partitioning, evolving schema/partitions, or fixing small-file and metadata bloat.

Unknown-333/building-ingestion-pipelines

Build batch and incremental data ingestion (extract-load) pipelines — full vs incremental extraction, change data capture (CDC), watermarks and high-water marks, API pagination and rate limits, and choosing managed EL tools (Fivetran, Airbyte) vs custom code. Use when ingesting data from databases, APIs, files, or SaaS into a warehouse/lake, or designing incremental extraction and CDC.

Unknown-333/debugging-data-pipelines

Systematically root-cause data pipeline failures and data incidents — job errors, wrong or missing data, duplicates, and freshness misses — by tracing lineage upstream, isolating the failing stage, reconciling against source, and planning a safe fix and backfill. Use when a pipeline fails, numbers look wrong, data is missing or duplicated, a dashboard is stale, or a stakeholder reports a data discrepancy.

Unknown-333/designing-backfills-and-replays

Plan and run safe data backfills and replays — idempotent reprocessing of historical windows, partition-by-partition execution, isolating backfill compute from production, verifying results, and avoiding double-counting or changed history. Use when backfilling a new or fixed model, reprocessing after a bug, replaying events, or loading history for a new pipeline without corrupting existing data.

Unknown-333/designing-data-contracts

Define and enforce data contracts between producers and consumers — explicit schema, semantics, ownership, SLAs, and versioning — to prevent silent upstream changes from breaking downstream pipelines. Use when a producer schema change could break consumers, defining an interface between teams/services and the warehouse, or adding schema enforcement at ingestion.

Verwandte Skills