with one click
data-engineering
数据工程(Airflow/Dagster/Kafka/Flink/dbt、数据管道、ETL、流处理、数据质量)。
Install with Codex or Claude Copy this prompt, paste it into Codex, Claude, or another assistant, and let it review the skill page and install it for you.
Menu
数据工程(Airflow/Dagster/Kafka/Flink/dbt、数据管道、ETL、流处理、数据质量)。
Install with Codex or Claude Copy this prompt, paste it into Codex, Claude, or another assistant, and let it review the skill page and install it for you.
Based on SOC occupation classification
CCG Skills - Quality gates, documentation generator, and multi-agent orchestration. Auto-installed by CCG workflow system.
变更校验关卡。分析代码变更,检测文档同步状态,评估变更影响范围。当用户提到变更检查、文档同步、代码审查、提交前检查、diff分析时使用。在设计级变更、重构完成时自动触发。
模块完整性校验关卡。扫描目录结构、检测缺失文档、验证代码与文档同步。当用户提到模块校验、文档检查、结构完整性、README检查、DESIGN检查时使用。在新建模块完成时自动触发。
代码质量校验关卡。检测复杂度、重复代码、命名规范、函数长度等质量指标。当用户提到代码质量、复杂度检查、代码异味、重构建议、lint检查、代码规范时使用。在复杂模块、重构完成时自动触发。
安全校验关卡。自动扫描代码安全漏洞,检测危险模式,确保安全决策有文档记录。当用户提到安全扫描、漏洞检测、安全审计、代码安全、OWASP、注入检测、敏感信息泄露时使用。在新建模块、安全相关变更、攻防任务、重构完成时自动触发。
从近期开发会话中萃取可复用的模式 / 决策 / 教训。当用户提到提取教训 / 复盘 / lessons learned / 知识沉淀 / 经验萃取 / 提炼模式时使用。会扫 .context/ 决策日志 + git log 提交语义。
| name | data-engineering |
| description | 数据工程(Airflow/Dagster/Kafka/Flink/dbt、数据管道、ETL、流处理、数据质量)。 |
| license | MIT |
| user-invocable | false |
| disable-model-invocation | false |
| context | fork |
数据工程域涵盖数据管道编排、流式处理、数据质量保障三大核心领域。
数据管道层 流处理层 质量保障层
├── Airflow (调度编排) ├── Kafka Streams ├── Great Expectations
├── Dagster (资产管理) ├── Flink ├── dbt
└── Prefect (现代工作流) └── Spark Streaming └── Soda Core
| 特性 | Airflow | Dagster | Prefect |
|---|---|---|---|
| 核心模型 | DAG + Task | Asset + Op | Flow + Task |
| 学习曲线 | 陡峭 | 中等 | 平缓 |
| 资产管理 | 无 | 原生支持 | 无 |
| 动态任务 | 支持 | 支持 | 支持 |
| 本地开发 | 复杂 | 简单 | 简单 |
| 社区生态 | 最大 | 成长中 | 成长中 |
with DAG(dag_id, schedule, default_args) as dag@task 装饰器,自动 XCom 传递@task + .expand() 实现 dynamic task mappingretries=3, retry_delay=timedelta(minutes=5), retry_exponential_backoff=Trueon_failure_callback 发送告警sla=timedelta(hours=2) + sla_miss_callback@asset(group_name, deps) 声明数据资产ConfigurableResource 管理外部连接define_asset_job(selection=AssetSelection.groups(...))ScheduleDefinition(job, cron_schedule)@sensor(job) 监听外部事件触发DailyPartitionsDefinition 按日分区@asset_check 验证数据新鲜度/质量@flow + @task(retries=3, cache_key_fn=task_input_hash)ConcurrentTaskRunner + task.map(items)Deployment.build_from_flow(schedule=CronSchedule(...))Secret / JSON 管理配置和密钥0 2 * * * 日批 / */15 * * * * 实时)WHERE updated_at > last_run| 特性 | Kafka Streams | Flink | Spark Streaming |
|---|---|---|---|
| 部署模式 | 嵌入式(JVM) | 独立集群 | 独立集群 |
| 状态管理 | RocksDB | 内存/RocksDB | 内存 |
| Exactly-Once | 支持 | 支持 | 支持 |
| 窗口类型 | 丰富 | 最丰富 | 基础 |
| 学习曲线 | 平缓 | 陡峭 | 中等 |
| Python API | kafka-python | PyFlink | PySpark |
StreamsBuilder → stream() → filter/map/flatMap → to()groupByKey().count() / .aggregate() / .reduce()Stores.persistentKeyValueStore + TransformerPROCESSING_GUARANTEE_CONFIG = EXACTLY_ONCE_V2NUM_STREAM_THREADS=4 / CACHE_MAX_BYTES_BUFFERING / RocksDB 配置env.addSource() → filter/map → addSink()TumblingProcessingTimeWindows.of(Time.minutes(5))SlidingProcessingTimeWindows.of(size, slide)ProcessingTimeSessionWindows.withGap(gap)GlobalWindows.create() + 自定义 Triggeraggregate(AggregateFunction, WindowFunction) 增量+全窗口env.enableCheckpointing(60000) + EXACTLY_ONCEflink run -s /path/to/savepointforBoundedOutOfOrderness)allowedLateness() + sideOutputLateData()完整性(非空) → 准确性(范围) → 一致性(关联) → 及时性(新鲜度) → 有效性(格式)
| 工具 | 优势 | 适用场景 |
|---|---|---|
| Great Expectations | 丰富 Expectations、Data Docs | Python 生态、复杂验证 |
| dbt | SQL 原生、血缘追踪 | 数据仓库、转换测试 |
| Soda Core | 简洁 YAML 配置 | 快速验证、CI/CD |
gx.get_context() → 添加数据源 → 构建批次expect_table_row_count_to_be_between(min, max)expect_column_values_to_not_be_null(column)expect_column_values_to_be_unique(column)expect_column_values_to_be_between(column, min, max)expect_column_values_to_be_in_set(column, value_set)expect_column_values_to_match_regex(column, regex)ColumnMapExpectationunique / not_null / accepted_values / relationships{% test name(model, column_name, params) %}tests/ 目录下自定义 SQL,返回行 = 失败expect_column_mean_to_be_between / expect_row_values_to_have_recent_datadbt test / dbt test --select model / dbt test --store-failures{{ ref('model') }} + {{ source('schema', 'table') }} → dbt docs generatechecks for table_name:
- row_count > 100
- missing_count(column) = 0
- duplicate_count(column) = 0
- invalid_count(column) = 0:
valid format: email
- freshness(timestamp_col) < 1d
| 实践 | 说明 |
|---|---|
| 幂等性设计 | UPSERT / 分区覆盖,重跑不产生副作用 |
| 增量处理 | 基于时间戳/CDC 增量提取,减少全量扫描 |
| 数据血缘 | dbt ref() / Dagster Asset deps 追踪上下游 |
| 分层验证 | 源→转换→目标每层都验证 |
| 监控告警 | 管道 SLA + 质量指标 + 延迟告警 |
| 状态管理 | 流处理状态 TTL + Checkpoint + Savepoint |
| 容错设计 | 重试策略 + 死信队列 + 回滚方案 |
数据管道、Airflow、Dagster、Prefect、ETL、流处理、Kafka Streams、Flink、数据质量、Great Expectations、dbt、数据验证、数据血缘