Contract
- Input: stream source descriptions, processing requirements, latency/throughput targets.
- Output: topology design + event schema + processing rules.
- Side effects: may deploy stream resources when executed.
- Dependencies: streaming platform access (Kafka / Kinesis / Pub/Sub / Flink).
- Stop condition: design saved; rules documented.
- Risk: medium — stream errors propagate quickly; requires testing.
- Boundary: designs stream; executes only with explicit approval.
Streaming Pipeline Design
Design a streaming pipeline — event producers, stream processors, consumers — with schema, processing rules, and reliability.
Process
1. Source analysis
- Event sources: user actions, sensors, logs, transactions, external APIs.
- Event rate (events/sec), burst rate, volume (GB/day).
- Event schema: fields, types, optional/required, nesting.
Completion criterion: source profile saved.
2. Event schema
- Define Avro / Protobuf / JSON Schema.
- Versioning: schema registry (Confluent Schema Registry / AWS Glue / GCP Pub/Sub schema).
- Backward/forward compatibility rules.
Completion criterion: schema saved; versioning rules defined.
3. Stream topology
- Producers: service A, B, C; publish to topic/stream.
- Stream: partitioned by key (e.g. user_id, device_id) for parallel processing.
- Processors: stream jobs (Flink / Kafka Streams / Spark Streaming / Kinesis Analytics) — filter, aggregate, enrich, transform.
- Consumers: service D reads results; may write to database / cache / warehouse.
Completion criterion: topology diagram saved.
4. Processing rules
- Windowing: tumbling, sliding, session, custom.
- Aggregation: count, sum, average, max, min; with watermarks.
- Enrichment: join with reference data (stream or database lookup).
- Exactly-once / at-least-once / at-most-once semantics per use case.
- Backpressure: when downstream is slow, how to handle (drop old, throttle, buffer with limit).
Completion criterion: rules saved.
5. Reliability
- Replication: stream replicas across zones.
- Retention: data retention period; archive to storage for long-term.
- Monitoring: lag per partition; error rate; throughput; processing latency.
- Alert: lag > threshold; error rate > threshold; partition unassigned.
Completion criterion: reliability rules saved.
1---2name: data-streaming3description: Design streaming data pipelines — Kafka, Kinesis, Pub/Sub, Flink — with event schemas, stream processing, and real-time analytics.4---56## Contract78- **Input:** stream source descriptions, processing requirements, latency/throughput targets.9- **Output:** topology design + event schema + processing rules.10- **Side effects:** may deploy stream resources when executed.11- **Dependencies:** streaming platform access (Kafka / Kinesis / Pub/Sub / Flink).12- **Stop condition:** design saved; rules documented.13- **Risk:** medium — stream errors propagate quickly; requires testing.14- **Boundary:** designs stream; executes only with explicit approval.1516# Streaming Pipeline Design1718Design a **streaming pipeline** — event producers, stream processors, consumers — with schema, processing rules, and reliability.1920## Process2122### 1. Source analysis23- Event sources: user actions, sensors, logs, transactions, external APIs.24- Event rate (events/sec), burst rate, volume (GB/day).25- Event schema: fields, types, optional/required, nesting.2627**Completion criterion:** source profile saved.2829### 2. Event schema30- Define Avro / Protobuf / JSON Schema.31- Versioning: schema registry (Confluent Schema Registry / AWS Glue / GCP Pub/Sub schema).32- Backward/forward compatibility rules.3334**Completion criterion:** schema saved; versioning rules defined.3536### 3. Stream topology37- **Producers:** service A, B, C; publish to topic/stream.38- **Stream:** partitioned by key (e.g. user_id, device_id) for parallel processing.39- **Processors:** stream jobs (Flink / Kafka Streams / Spark Streaming / Kinesis Analytics) — filter, aggregate, enrich, transform.40- **Consumers:** service D reads results; may write to database / cache / warehouse.4142**Completion criterion:** topology diagram saved.4344### 4. Processing rules45- Windowing: tumbling, sliding, session, custom.46- Aggregation: count, sum, average, max, min; with watermarks.47- Enrichment: join with reference data (stream or database lookup).48- Exactly-once / at-least-once / at-most-once semantics per use case.49- Backpressure: when downstream is slow, how to handle (drop old, throttle, buffer with limit).5051**Completion criterion:** rules saved.5253### 5. Reliability54- Replication: stream replicas across zones.55- Retention: data retention period; archive to storage for long-term.56- Monitoring: lag per partition; error rate; throughput; processing latency.57- Alert: lag > threshold; error rate > threshold; partition unassigned.5859**Completion criterion:** reliability rules saved.