ワンクリックで
data-engineering
数据工程(Airflow/Dagster/Kafka/Flink/dbt、数据管道、ETL、流处理、数据质量)。
Codex または Claude でインストール この Prompt をコピーして Codex、Claude、または他のアシスタントに貼り付けると、Skill ページを確認してインストールできます。
メニュー
数据工程(Airflow/Dagster/Kafka/Flink/dbt、数据管道、ETL、流处理、数据质量)。
Codex または Claude でインストール この Prompt をコピーして Codex、Claude、または他のアシスタントに貼り付けると、Skill ページを確認してインストールできます。
SOC 職業分類に基づく
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、数据验证、数据血缘