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

如何在Airflow的PostgresOperator中记录增改删行数并可视化?

实现SQL操作行数记录与可视化方案

一、按SQL文件分别记录操作行数

1. 拆分独立SQL任务

当前你把多个SQL文件放在单个PostgresOperator中执行,无法单独追踪每个文件的执行结果。建议将每个SQL拆分为独立任务,便于精准记录行数:

from airflow import DAG
from airflow.providers.postgres.operators.postgres import PostgresOperator
from datetime import datetime

with DAG(
    dag_id="postgres_operator_dag",
    start_date=datetime(2023, 2, 2),
    schedule_interval=None,
    catchup=False,
) as dag:
    task_001 = PostgresOperator(
        task_id='execute_001_test',
        postgres_conn_id='postgres_dbad2a',
        sql='001-test.sql'
    )
    task_002 = PostgresOperator(
        task_id='execute_002_test',
        postgres_conn_id='postgres_dbad2a',
        sql='002-test.sql'
    )

    task_001 >> task_002

2. 捕获DML行数的两种方案

方案A:修改SQL嵌入行数捕获逻辑

在每个SQL文件的DML语句后添加PostgreSQL原生的行数捕获代码,执行后日志会自动输出行数:

-- 001-test.sql示例
INSERT INTO your_table (col1, col2) VALUES ('val1', 'val2'), ('val3', 'val4');
GET DIAGNOSTICS inserted_rows = ROW_COUNT;
RAISE NOTICE 'INSERTED ROWS: %', inserted_rows;

UPDATE your_table SET col2 = 'updated' WHERE col1 = 'val1';
GET DIAGNOSTICS updated_rows = ROW_COUNT;
RAISE NOTICE 'UPDATED ROWS: %', updated_rows;

任务执行后,可在Airflow日志中直接查看每个DML的受影响行数。

方案B:自定义Operator存储行数到XCom

继承PostgresOperator,自动执行SQL并捕获行数,通过XCom存储数据(方便后续可视化):

from airflow.providers.postgres.operators.postgres import PostgresOperator
from airflow.providers.postgres.hooks.postgres import PostgresHook
from airflow.utils.decorators import apply_defaults

class PostgresDmlCountOperator(PostgresOperator):
    @apply_defaults
    def __init__(self, **kwargs):
        super().__init__(**kwargs)

    def execute(self, context):
        hook = PostgresHook(postgres_conn_id=self.postgres_conn_id)
        conn = hook.get_conn()
        cursor = conn.cursor()
        
        # 读取SQL文件内容
        if isinstance(self.sql, str) and self.sql.endswith('.sql'):
            with open(self.sql, 'r') as f:
                sql_content = f.read()
        else:
            sql_content = self.sql
        
        # 拆分SQL语句(按分号分割)
        sql_statements = [stmt.strip() for stmt in sql_content.split(';') if stmt.strip()]
        
        row_counts = {}
        for idx, stmt in enumerate(sql_statements):
            cursor.execute(stmt)
            row_count = cursor.rowcount
            stmt_type = stmt.split()[0].upper()
            if stmt_type in ['INSERT', 'UPDATE', 'DELETE']:
                row_counts[f"{stmt_type.lower()}_rows_{idx+1}"] = row_count
        
        # 将行数存入XCom
        context['ti'].xcom_push(key='dml_row_counts', value=row_counts)
        
        conn.commit()
        cursor.close()
        conn.close()

使用示例:

with DAG(...) as dag:
    task_001 = PostgresDmlCountOperator(
        task_id='execute_001_test',
        postgres_conn_id='postgres_dbad2a',
        sql='001-test.sql'
    )

执行后可在任务的XCom面板查看存储的行数数据。

二、行数数据的可视化展示

1. 自定义Airflow UI插件

开发Airflow UI插件,读取XCom中的行数数据并渲染图表:

  • 从Airflow元数据库的xcom表中查询指定DAG、任务的dml_row_counts数据
  • 用ECharts、Chart.js等前端库生成柱状图/折线图,嵌入Airflow UI

2. Grafana可视化方案

通过Grafana连接Airflow元数据库,直接查询并展示行数:

  • 配置Grafana连接Airflow的元数据库(如PostgreSQL)
  • 编写查询SQL提取XCom数据:
SELECT
    ti.dag_id,
    ti.task_id,
    (xcom.value::json->>'inserted_rows')::int AS inserted_rows,
    (xcom.value::json->>'updated_rows')::int AS updated_rows,
    (xcom.value::json->>'deleted_rows')::int AS deleted_rows,
    ti.start_date
FROM task_instance ti
JOIN xcom ON ti.task_id = xcom.task_id AND ti.dag_id = xcom.dag_id
WHERE xcom.key = 'dml_row_counts'
ORDER BY ti.start_date DESC
  • 选择柱状图、折线图等类型,创建自定义仪表盘

3. Prometheus+Grafana指标监控

将行数作为自定义指标推送至Prometheus,再通过Grafana展示:

  • 在自定义Operator中添加指标上报逻辑:
from airflow.metrics import Metrics

# 在execute方法中添加
for metric_key, count in row_counts.items():
    Metrics.gauge(
        f"airflow_dml_{metric_key}", 
        count, 
        tags={"dag_id": context['dag'].dag_id, "task_id": self.task_id}
    )
  • 配置Airflow将指标推送至Prometheus,在Grafana中创建监控面板

内容的提问来源于stack exchange,提问作者Rob Audenaerde

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 08:08:52