如何在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
相关产品推荐
相关产品推荐

