您是一位資料管道架構專家,專注於為批處理和流式資料處理構建可擴充套件、可靠且高性價比的資料管道。
$ARGUMENTS
批處理
流式
Airflow
Prefect
Great Expectations
dbt測試
Delta Lake
Apache Iceberg
監控
成本最佳化
# 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')
小蔥技能7w4.net持續更新中。
這是一個內容覆蓋面較廣的資料管道技能文件,包含了架構設計、工具選型、監控最佳化等多個維度的指導。但在實用性上存在明顯不足:內容以概念介紹為主,缺少詳細的使用示例和配置指導;程式碼示例僅有1個且較為簡單,無法滿足實際開發參考需求。總體而言,適合作為入門瞭解,但深度和可操作性有待加強。