You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Airflow operator中调用BingAds API获取营销活动数据?

基于Airflow实现Bing Ads数据下载到S3的落地方案

前期依赖准备

  • 先安装必要依赖:bing-ads Python SDK、apache-airflow-providers-amazon(用于S3 Hook),两个依赖包都需要预安装到Airflow worker的运行环境中,避免运行时报依赖缺失错误。
  • 提前在Bing Ads后台申请开发者token、OAuth2凭据(client_id、client_secret、refresh_token),以及对应广告账户的customer_id、account_id,所有敏感信息统一存入Airflow Connections或Variables中,禁止硬编码到代码里。

核心实现逻辑

你直接用PythonOperator封装逻辑完全可行,不需要额外开发自定义Operator,核心执行步骤如下:

  1. 在Python回调函数中初始化Bing Ads SDK客户端,从Airflow配置项拉取认证信息完成鉴权。
  2. 调用Bing Ads的报表服务接口,按需传入要拉取的营销活动维度、指标参数,SDK原生支持直接返回CSV格式结果,不需要手动把Python对象列表转成文件。
  3. 如果有数据清洗需求可以在内存中处理完再写入临时文件,也可以直接把SDK返回的文件流通过S3 Hook上传到指定S3路径,无需本地落盘减少IO开销。
  4. 上传完成后清理本地临时文件(如有落盘),返回任务执行成功状态即可。

简化版代码示例

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.amazon.aws.hooks.s3 import S3Hook
from airflow.models import Variable
from bingads import AuthorizationData, ServiceClient
from bingads.v13.reporting import *
from datetime import datetime
import tempfile
import os

def pull_bing_ads_data_to_s3(**context):
    # 从Airflow配置拉取认证信息,生产环境建议用Airflow Connection存储更安全
    auth_data = AuthorizationData(
        developer_token=Variable.get("bing_ads_developer_token"),
        client_id=Variable.get("bing_ads_client_id"),
        client_secret=Variable.get("bing_ads_client_secret"),
        refresh_token=Variable.get("bing_ads_refresh_token"),
        account_id=Variable.get("bing_ads_account_id"),
        customer_id=Variable.get("bing_ads_customer_id")
    )

    # 初始化报表服务客户端
    reporting_service = ServiceClient(
        service='ReportingService',
        version=13,
        authorization_data=auth_data,
        environment='production'
    )

    # 构造报表请求,此处以CampaignPerformanceReport为例,可按需替换为其他报表类型
    report_request = reporting_service.factory.create('CampaignPerformanceReportRequest')
    report_request.Format = 'Csv'
    report_request.ReportName = 'DailyCampaignPerformanceReport'
    report_request.ReturnOnlyCompleteData = True
    # 配置需要的报表字段
    report_request.Columns = reporting_service.factory.create('ArrayOfCampaignPerformanceReportColumn')
    report_request.Columns.CampaignPerformanceReportColumn.append([
        'TimePeriod', 'CampaignName', 'CampaignId', 'Impressions', 'Clicks', 'Spend', 'Conversions'
    ])
    # 配置拉取时间范围,适配Airflow调度的逻辑日期
    report_time = reporting_service.factory.create('ReportTime')
    report_time.PredefinedTime = 'Yesterday'
    report_request.Time = report_time

    # 提交请求下载报表到临时目录
    reporting_download_parameters = ReportingDownloadParameters(
        report_request=report_request,
        result_file_directory=tempfile.gettempdir(),
        result_file_name=f"bing_ads_campaign_{context['ds']}.csv",
        overwrite_result_file=True
    )
    file_path = ReportingServiceManager(authorization_data=auth_data).download_file(reporting_download_parameters)

    # 上传文件到S3
    s3_hook = S3Hook(aws_conn_id='your_aws_conn_id')
    s3_hook.load_file(
        filename=file_path,
        key=f"bing_ads/campaign/dt={context['ds']}/data.csv",
        bucket_name='your_s3_bucket_name',
        replace=True
    )

    # 清理本地临时文件
    os.remove(file_path)

# DAG基础定义部分可根据你的实际调度规则自行配置即可

注意事项

  • Bing Ads API有并发和频率限制,同一开发者token的请求不要过于密集,任务可以设置2-3次重试,重试间隔设为5分钟以上。
  • refresh_token存在有效期,建议单独配置一个周级调度的小任务定期刷新token并更新到Airflow配置中,避免认证失败。
  • 如果需要拉取大跨度历史数据,建议按天拆分请求分批拉取,避免单个报表请求超时或者返回数据量过大。

内容的提问来源于stack exchange,提问作者Ramon Soto

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.24 19:15:03