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

如何在DBT中创建用于外部API调用和JDBC连接的Python宏

在DBT中构建Python宏实现Azure Monitor API调用与JDBC集成

一、DBT Python宏基础结构

DBT 1.3+版本支持Python宏,直接在项目的macros/目录下创建.py文件(比如azure_integration.py),函数即为可调用的宏。无需额外装饰器,DBT会自动识别该目录下的Python文件。

先搭建基础框架,引入必要依赖:

# macros/azure_integration.py
import requests
import jaydebeapi
from dbt import context
from dbt.logger import logger
from azure.identity import DefaultAzureCredential

二、Azure Monitor API调用实现

Azure Monitor API需通过Azure AD认证,先实现Token获取逻辑,再封装API调用:

def get_azure_ad_token(resource="https://monitor.azure.com/"):
    """获取Azure AD认证Token"""
    credential = DefaultAzureCredential()
    token = credential.get_token(resource)
    return token.token

def fetch_azure_monitor_metrics(workspace_id, metric_name, timespan):
    """调用Azure Monitor API获取指标数据"""
    # 从DBT变量中读取全局配置
    config = context.get('vars')
    api_url = (
        f"https://api.monitor.azure.com/subscriptions/{config['azure_subscription_id']}"
        f"/resourcegroups/{config['azure_resource_group']}"
        f"/providers/microsoft.insights/components/{workspace_id}/metrics"
    )

    token = get_azure_ad_token()
    headers = {
        "Authorization": f"Bearer {token}",
        "Content-Type": "application/json"
    }
    params = {
        "metricnames": metric_name,
        "timespan": timespan,
        "aggregation": "average"
    }

    logger.info(f"发起Azure Monitor API请求: {api_url}")
    response = requests.get(api_url, headers=headers, params=params)
    response.raise_for_status()
    return response.json()

三、JDBC连接与数据交互

使用jaydebeapi库实现JDBC连接,封装查询、插入逻辑:

def connect_jdbc(jdbc_url, driver_class, driver_path, credentials):
    """建立JDBC连接"""
    logger.info(f"尝试连接JDBC数据源: {jdbc_url}")
    conn = jaydebeapi.connect(
        driver_class,
        jdbc_url,
        credentials,
        driver_path
    )
    return conn

def insert_api_results_to_jdbc(conn, api_data):
    """将API返回数据插入JDBC数据源"""
    insert_query = """
        INSERT INTO monitor_metrics (metric_name, metric_value, record_timestamp)
        VALUES (?, ?, ?)
    """
    cursor = conn.cursor()
    metric_count = 0

    for item in api_data.get('value', []):
        for data_point in item['timeseries'][0]['data']:
            cursor.execute(
                insert_query,
                (
                    item['name']['value'],
                    data_point.get('average'),
                    data_point['timeStamp']
                )
            )
            metric_count += 1
    
    conn.commit()
    cursor.close()
    logger.info(f"成功插入{metric_count}条监控指标数据")

四、整合逻辑并集成到DBT工作流

将API调用与JDBC操作整合成一个可直接触发的宏,再通过多种方式集成到工作流:

1. 封装整合宏

def sync_azure_monitor_to_jdbc():
    """端到端同步Azure Monitor数据到JDBC数据源"""
    config = context.get('vars')
    jdbc_conn = None

    try:
        # 1. 获取API数据
        metrics_data = fetch_azure_monitor_metrics(
            workspace_id=config['azure_monitor_workspace_id'],
            metric_name=config['metric_name'],
            timespan=config['timespan']
        )

        # 2. 建立JDBC连接
        jdbc_conn = connect_jdbc(
            jdbc_url=config['jdbc_url'],
            driver_class=config['jdbc_driver_class'],
            driver_path=config['jdbc_driver_path'],
            credentials={
                "user": config['jdbc_user'],
                "password": config['jdbc_password']
            }
        )

        # 3. 插入数据
        insert_api_results_to_jdbc(jdbc_conn, metrics_data)

    except requests.exceptions.HTTPError as e:
        raise Exception(f"Azure Monitor API调用失败: {str(e)}")
    except jaydebeapi.DatabaseError as e:
        raise Exception(f"JDBC操作失败: {str(e)}")
    finally:
        if jdbc_conn:
            jdbc_conn.close()
            logger.info("JDBC连接已关闭")

2. 触发方式

  • 在模型中调用:创建一个临时模型触发宏执行(仅用于触发逻辑,无实际数据输出)
-- models/staging/sync_monitor_data.sql
{{ config(materialized='ephemeral') }}

{% do sync_azure_monitor_to_jdbc() %}

SELECT 1 AS dummy;

运行dbt run --model sync_monitor_data即可触发同步。

  • 直接用命令调用:通过dbt run-operation直接执行宏
dbt run-operation sync_azure_monitor_to_jdbc --vars '{
    "azure_subscription_id": "你的订阅ID",
    "azure_resource_group": "你的资源组",
    "azure_monitor_workspace_id": "你的Monitor工作区ID",
    "metric_name": "Requests",
    "timespan": "PT1H",
    "jdbc_url": "jdbc:sqlserver://xxx.database.windows.net:1433;database=xxx;",
    "jdbc_driver_class": "com.microsoft.sqlserver.jdbc.SQLServerDriver",
    "jdbc_driver_path": "./drivers/mssql-jdbc-12.4.1.jre11.jar",
    "jdbc_user": "你的数据库用户",
    "jdbc_password": "你的数据库密码"
}'
  • 集成到DBT作业:在DBT Cloud或本地作业中,将上述dbt run-operation命令作为步骤之一,与其他模型运行逻辑串联。

五、最佳实践

  • 凭据安全管理:避免硬编码敏感信息,通过DBT变量结合环境变量传递
# dbt_project.yml
vars:
  azure_subscription_id: "{{ env_var('AZURE_SUBSCRIPTION_ID') }}"
  jdbc_password: "{{ env_var('JDBC_PASSWORD') }}"

运行时通过环境变量注入敏感值。

  • 依赖管理:在项目根目录创建requirements.txt,列出所有依赖包
requests>=2.31.0
jaydebeapi>=1.2.3
azure-identity>=1.12.0

执行pip install -r requirements.txt安装依赖,DBT Cloud可直接识别该文件自动安装。

  • 日志与错误处理:通过dbt.logger记录关键步骤,添加异常捕获确保工作流稳定性。

内容的提问来源于stack exchange,提问作者Ankit Devani

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 19:14:57