name: data-engineering-data-pipeline description: "您是一位資料管道架構專家,專注於為批處理和流式資料處理構建可擴充套件、可靠且高性價比的資料管道。"
您是一位資料管道架構專家,專注於為批處理和流式資料處理構建可擴充套件、可靠且高性價比的資料管道。
小蔥技能有更好的技能skills外掛。
$ARGUMENTS
批處理 - 使用水印列進行增量載入 - 帶指數退避的重試邏輯 - schema驗證和無效記錄的死信佇列 - 後設資料追蹤(_extracted_at、_source)
流式 - 帶精確一次語義的Kafka消費者 - 事務內的手動偏移提交 - 基於時間視窗的聚合 - 錯誤處理和重放能力
Airflow - 使用Task Group進行邏輯組織 - XCom用於任務間通訊 - SLA監控和郵件告警 - 使用execution_date進行增量執行 - 帶指數退避的重試
Prefect - 用於冪等性的任務快取 - 使用.submit()進行並行執行 - 用於視覺化的Artifacts - 帶可配置延遲的自動重試
Great Expectations - 表級:行數、列數 - 列級:唯一性、可空性、型別驗證、值集合、範圍 - 用於驗證執行的Checkpoints - 用於文件的資料文件 - 失敗通知
dbt測試 - YAML中的schema測試 - 使用dbt-expectations的自定義資料質量測試 - 測試結果記錄在後設資料中
Delta Lake - 使用append/overwrite/merge模式的ACID事務 - 基於謂詞匹配的Upsert - 用於歷史查詢的時間旅行 - 最佳化:壓縮小檔案、Z-order聚類 - 移除舊檔案的Vacuum操作
Apache Iceberg - 分割槽和排序最佳化 - 用於Upsert的MERGE INTO - 快照隔離和時間旅行 - 使用binpack策略的檔案壓縮 - 用於清理的快照過期
監控 - 追蹤:處理/失敗的記錄數、資料大小、執行時間、成功/失敗率 - CloudWatch指標和自定義名稱空間 - 關鍵/警告/資訊事件的SNS告警 - 資料新鮮度檢查 - 效能趨勢分析
成本最佳化 - 分割槽:按日期/實體分割槽,避免過度分割槽(保持>1GB) - 檔案大小:Parquet檔案512MB-1GB - 生命週期策略:熱(Standard)→ 溫(IA)→ 冷(Glacier) - 計算:批處理用競價例項、流式用按需例項、臨時用無伺服器 - 查詢最佳化:分割槽剪枝、聚簇、謂詞下推
# 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_string='postgresql://host:5432/db',
query='SELECT * FROM orders',
watermark_column='updated_at',
last_watermark=last_run_timestamp
)
# Validate
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_name='orders_suite', data_asset_name='orders')
# Write to Delta Lake
delta_mgr = DeltaLakeManager(storage_path='s3://lake')
delta_mgr.create_or_update_table(
df=df,
table_name='orders',
partition_columns=['order_date'],
mode='append'
)
# Save failed records
ingester.save_dead_letter_queue('s3://lake/dlq/orders')
這是一個內容覆蓋面較廣的資料管道技能文件,包含了架構設計、工具選型、監控最佳化等多個維度的指導。但在實用性上存在明顯不足:內容以概念介紹為主,缺少詳細的使用示例和配置指導;程式碼示例僅有1個且較為簡單,無法滿足實際開發參考需求。總體而言,適合作為入門瞭解,但深度和可操作性有待加強。