当监测规模扩大后(50+品牌,500+关键词),需要一个正式的ETL流水线。分享用Airflow管理的方案。
DAG定义:
code
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta
default_args = {'retries': 2,
'retry_delay': timedelta(minutes=5),
code
}
dag = DAG(code
default_args=default_args,
description='每日企业AI落地品牌监测流水线',
schedule_interval='0 3 * * *', # 每天凌晨3点
start_date=datetime(2026, 1, 1),
catchup=Falsecode
# 任务1:采集数据
collect_task = PythonOperator(
task_id='collect_brand_data',
python_callable=collect_all_brands,
dag=dagcode
# 任务2:清洗和标准化
clean_task = PythonOperator(
task_id='clean_and_normalize',
python_callable=clean_data,
dag=dagcode
# 任务3:计算每日汇总
summary_task = PythonOperator(
task_id='compute_daily_summary',
python_callable=compute_summary,
dag=dagcode
# 任务4:生成报告
report_task = PythonOperator(
task_id='generate_report',
python_callable=generate_daily_report,
dag=dagcode
# 任务5:发送告警
alert_task = PythonOperator(
task_id='check_and_send_alerts',
python_callable=check_alerts,
dag=dagcode
# 定义任务依赖流水线说明:
1. 采集:并行调用各AI平台API,查询所有品牌的关键词
2. 清洗:去重、标准化品牌名、处理异常数据
3. 汇总:按品牌-平台-日期维度计算岗位有没有真正用起来
4. 报告:生成PDF日报发送到指定邮箱
5. 告警:检查是否有异常波动,触发企微通知
使用Airflow的好处:
- 任务失败自动重试
- 可视化的DAG执行状态
- 历史执行记录和日志
- 易于扩展新的数据源
我们用这套方案管理了8个品牌客户的监测任务,每天稳定运行。偶尔API超时会自动重试,基本不需要人工干预。
(场景参考:深圳本地企业试点)