Loading...
Loading...
Compare original and translation side by side
dbt-skilldbt-skillreferences/kafka-deep-dive.mdreferences/flink-spark-streaming.mdreferences/warehouse-streaming-ingestion.mdreferences/stream-testing-patterns.mdreferences/kafka-deep-dive.mdreferences/flink-spark-streaming.mdreferences/warehouse-streaming-ingestion.mdreferences/stream-testing-patterns.md| reasoning_demand | preferred | acceptable | minimum |
|---|---|---|---|
| high | Opus | Sonnet | Sonnet |
| 推理需求 | 首选 | 可接受 | 最低要求 |
|---|---|---|---|
| 高 | Opus | Sonnet | Sonnet |
WATERMARK FOR event_timestamp AS event_timestamp - INTERVAL 'N' SECONDWATERMARK FOR event_timestamp AS event_timestamp - INTERVAL 'N' SECOND| Architecture | Latency | Complexity | Best For |
|---|---|---|---|
| Traditional Batch | Hours-days | Low | Historical reporting, large aggregations |
| Micro-Batch (Spark Streaming) | Seconds-minutes | Medium | Near-real-time analytics, unified batch/stream |
| True Streaming (Flink, Kafka Streams) | Milliseconds-seconds | High | Real-time dashboards, fraud detection, alerting |
| Kappa Architecture | Milliseconds-seconds | Medium | Stream-first, immutable event log, reprocessing |
| Warehouse-Native Streaming | Seconds-minutes | Low | SQL-first teams, simple ingestion, BI integration |
| 架构 | 延迟 | 复杂度 | 最佳适用场景 |
|---|---|---|---|
| 传统批处理 | 数小时-数天 | 低 | 历史报告、大规模聚合 |
| 微批处理(Spark Streaming) | 数秒-数分钟 | 中 | 近实时分析、批流统一 |
| 真正流处理(Flink、Kafka Streams) | 毫秒-数秒 | 高 | 实时仪表板、欺诈检测、告警 |
| Kappa架构 | 毫秒-数秒 | 中 | 流优先、不可变事件日志、重处理 |
| 数据仓库原生流处理 | 数秒-数分钟 | 低 | 以SQL为主的团队、简单摄入、BI集成 |
| Framework | Latency | SQL | Managed | Best For |
|---|---|---|---|---|
| Kafka Streams | ms | KSQL (separate) | No | Microservices, JVM apps |
| Apache Flink | ms | Flink SQL | AWS KDA, Confluent | Complex event processing, large state |
| Spark Structured Streaming | seconds | Spark SQL | Databricks, EMR | Unified batch/stream, ML integration |
| ksqlDB | ms | Streaming SQL | Confluent Cloud | SQL-first simple transforms |
| Apache Beam/Dataflow | seconds | Limited | GCP Dataflow | Multi-cloud, GCP native |
| 框架 | 延迟 | SQL支持 | 托管服务 | 最佳适用场景 |
|---|---|---|---|---|
| Kafka Streams | 毫秒 | KSQL(独立) | 无 | 微服务、JVM应用 |
| Apache Flink | 毫秒 | Flink SQL | AWS KDA、Confluent | 复杂事件处理、大规模状态 |
| Spark Structured Streaming | 数秒 | Spark SQL | Databricks、EMR | 批流统一、机器学习集成 |
| ksqlDB | 毫秒 | 流SQL | Confluent Cloud | 以SQL为主的简单转换 |
| Apache Beam/Dataflow | 数秒 | 有限 | GCP Dataflow | 多云、GCP原生 |
| Pattern | Description | Use Case |
|---|---|---|
| Tumbling | Fixed, non-overlapping intervals | Hourly aggregations, regular reporting |
| Sliding | Fixed, overlapping intervals | Moving averages, trend detection |
| Session | Gap-based, variable size | User sessions, activity bursts |
| Global | Custom trigger-controlled | Accumulate until condition met |
| 模式 | 描述 | 适用场景 |
|---|---|---|
| Tumbling(滚动窗口) | 固定、无重叠的时间间隔 | 小时级聚合、常规报告 |
| Sliding(滑动窗口) | 固定、有重叠的时间间隔 | 移动平均值、趋势检测 |
| Session(会话窗口) | 基于间隙的可变大小窗口 | 用户会话、活动爆发 |
| Global(全局窗口) | 自定义触发器控制的窗口 | 累积数据直到满足条件 |
| Mode | Allowed Changes | Use When |
|---|---|---|
| BACKWARD | Delete fields, add optional | Consumers upgrade first (most common) |
| FORWARD | Add fields, delete optional | Producers upgrade first |
| FULL | Backward + Forward only | Upgrade order unpredictable |
| 模式 | 允许的变更 | 使用场景 |
|---|---|---|
| BACKWARD | 删除字段、添加可选字段 | 先升级消费者(最常见) |
| FORWARD | 添加字段、删除可选字段 | 先升级生产者 |
| FULL | 仅允许BACKWARD + FORWARD变更 | 升级顺序不可预测 |
| Metric | Alert Threshold |
|---|---|
| Consumer Lag | > 1M messages or > 5 min |
| Throughput | < 50% baseline |
| Error Rate | > 0.1% for critical pipelines |
| Checkpoint Duration | > 2x interval |
| Backpressure Ratio | > 10% sustained |
| Partition Skew | Max/min ratio > 3x |
| 指标 | 告警阈值 |
|---|---|
| 消费者延迟 | > 100万条消息 或 > 5分钟 |
| 吞吐量 | < 基准值的50% |
| 错误率 | 关键管道 > 0.1% |
| 检查点持续时间 | > 间隔的2倍 |
| 背压比率 | 持续 > 10% |
| 分区倾斜 | 最大/最小比率 > 3倍 |
ConfigProvider| Capability | Tier 1 (Cloud-Native) | Tier 2 (Regulated) | Tier 3 (Air-Gapped) |
|---|---|---|---|
| Kafka producer/consumer | Deploy to dev | Generate for review | Generate only |
| Flink/Spark jobs | Submit to dev | Generate for review | Generate only |
| Warehouse streaming | Configure dev | Generate configs | Generate only |
ConfigProvider| 能力 | Tier 1(云原生) | Tier 2(受监管) | Tier 3(离线环境) |
|---|---|---|---|
| Kafka生产者/消费者 | 部署到开发环境 | 生成供审核 | 仅生成 |
| Flink/Spark作业 | 提交到开发环境 | 生成供审核 | 仅生成 |
| 数据仓库流处理 | 配置开发环境 | 生成配置 | 仅生成 |