Data Pipeline 数据管道架构设计实战从 ETL/ELT 选型到 Lakehouse 存储与质量监控全链路【免费下载链接】agentsMulti-harness agentic plugin marketplace for Claude Code, Codex, Cursor, OpenCode, GitHub Copilot, and Google Antigravity项目地址: https://gitcode.com/GitHub_Trending/agents24/agents本篇文章围绕 agents24 仓库中># dags/example_dag.py from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator from airflow.operators.empty import EmptyOperator default_args { owner: data-team, depends_on_past: False, email_on_failure: True, email_on_retry: False, retries: 3, retry_delay: timedelta(minutes5), retry_exponential_backoff: True, max_retry_delay: timedelta(hours1), } with DAG( dag_idexample_etl, default_argsdefault_args, descriptionExample ETL pipeline, schedule0 6 * * *, # Daily at 6 AM start_datedatetime(2024, 1, 1), catchupFalse, tags[etl, example], max_active_runs1, ) as dag: start EmptyOperator(task_idstart) def extract_data(**context): execution_date context[ds] # 使用宏而非硬编码日期 return {records: 1000} extract PythonOperator(task_idextract, python_callableextract_data) end EmptyOperator(task_idend) start extract end该技能的 references/details.md 还覆盖了 TaskFlow API、动态 DAG 工厂、分支逻辑BranchPythonOperator TriggerRule、S3KeySensor/ExternalTaskSensor 传感器、失败/成功回调告警以及 pytest 驱动的 DAG 结构测试形成写 DAG → 测 DAG → 报警的完整闭环。4.2 Prefect 四大实践任务缓存以缓存保证幂等避免重复计算并行执行通过.submit()提交并行任务Artifacts 可见性把关键产出物数据量、质量分数暴露为可观测工件自动重试与可配置延迟按任务维度配置重试次数与间隔。五、dbt 转换staging → marts 的分层纪律命令对 dbt 转换层提出五条要求Staging 层增量物化incremental materialization、去重、晚到数据处理late-arriving dataMarts 层维度模型、聚合、业务逻辑落地测试unique、not_null、relationships、accepted_values及自定义数据质量测试Sources新鲜度检查freshness checks、loaded_at_field追踪增量策略merge或deleteinsert两种主流方案。仓库 dbt-transformation-patterns 提供了可直接落地的dbt_project.yml与分层目录结构。其中增量策略的实现细节值得展开见 references/details.md-- deleteinsert多数数仓默认 {{ config( materializedincremental, unique_keyid, incremental_strategydeleteinsert ) }} -- merge适合晚到数据修正 {{ config( materializedincremental, unique_keyid, incremental_strategymerge, merge_update_columns[status, amount, updated_at] ) }} -- insert_overwrite基于分区的覆盖写入 {{ config( materializedincremental, incremental_strategyinsert_overwrite, partition_by{ field: created_date, data_type: date, granularity: day } ) }}配套命令集还给出开发期常用命令dbt run、dbt run --select fct_orders含上游、dbt run --full-refresh、dbt test、dbt build按 DAG 顺序 run test、dbt docs generate/dbt docs serve、dbt compile/dbt debug。六、数据质量框架Great Expectations 与 dbt Tests 双轨并行6.1 Great Expectations命令要求覆盖表级与列级两个维度表级行数、列数如expect_table_row_count_to_be_between列级唯一性、空值率、类型校验、取值集合、数值范围如expect_column_values_to_be_unique、expect_column_values_to_be_in_set、expect_column_values_to_be_betweenCheckpoints定义校验的执行载体与调度入口Data Docs自动生成可读的数据质量文档失败通知校验失败时触发告警链路。仓库># Batch ingestion with validation from batch_ingestion import BatchDataIngester from storage.delta_lake_manager import DeltaLakeManager from data_quality.expectations_suite import DataQualityFramework ingester BatchDataIngester(config{}) # Extract with incremental loading —— 水位列增量抽取 df ingester.extract_from_database( connection_stringpostgresql://host:5432/db, querySELECT * FROM orders, watermark_columnupdated_at, last_watermarklast_run_timestamp ) # Validate —— 摄取后立即做 schema 校验 schema {required_fields: [id, user_id], dtypes: {id: int64}} df ingester.validate_and_clean(df, schema) # Data quality checks —— 进入存储前过质量门 dq DataQualityFramework() result dq.validate_dataframe(df, suite_nameorders_suite, data_asset_nameorders) # Write to Delta Lake —— 分区追加写入 delta_mgr DeltaLakeManager(storage_paths3://lake) delta_mgr.create_or_update_table( dfdf, table_nameorders, partition_columns[order_date], modeappend ) # Save failed records —— 坏数据进死信队列不阻塞主链路 ingester.save_dead_letter_queue(s3://lake/dlq/orders)该示例完整演绎了增量抽取 → schema 校验 → 质量检查 → 分区写入 → 死信队列五步流水线将本文第三、六、七节的原则固化为可读的代码形态。注意示例中的类名BatchDataIngester、DeltaLakeManager、DataQualityFramework是命令给出的接口示意实际落地时应替换为对应生态的真实实现。十、输出交付物与验收标准10.1 五类交付物命令要求管道专家按五类可审查的产物交付而非只交代码架构文档数据流架构图、技术栈选型及理由、可扩展性分析、故障模式与恢复策略实现代码摄取批/流 错误处理、转换dbt models 或 Spark jobs、编排Airflow/Prefect DAGs、存储Delta/Iceberg 表管理、质量GE suites 与 dbt tests配置文件编排 DAG 定义/调度/重试策略、dbt 的 models/sources/tests/project config、基础设施Docker Compose、K8s manifests、Terraform、环境配置dev/staging/prod监控与可观测性指标执行时间、处理记录数、质量分数、告警失败、性能劣化、数据新鲜度、仪表盘Grafana/CloudWatch、结构化日志correlation IDs运维指南部署与回滚流程、常见问题排障手册、扩容指南、成本优化策略、灾备与备份流程。10.2 七项成功标准管道满足既定 SLA延迟、吞吐数据质量检查通过率 99%失败自动重试与告警全面监控展示健康度与性能文档足以支撑团队独立维护成本优化使基础设施成本下降 30–50%无停机 Schema 演化端到端数据血缘可追踪。十一、如何在仓库中使用该命令本命令属于plugins/data-engineering插件其配套资源在仓库中的组织方式如下命令定义data-pipeline.md 与 contenteditable="false">【免费下载链接】agentsMulti-harness agentic plugin marketplace for Claude Code, Codex, Cursor, OpenCode, GitHub Copilot, and Google Antigravity项目地址: https://gitcode.com/GitHub_Trending/agents24/agents创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考