基于湖库一体架构,统一管理结构化、半结构化与非结构化等多模态数据,一个系统承载事务处理、实时分析与 AI 工作负载。
Airflow 集成 OceanBase 数据库
更新时间:2026-05-18 16:30:49
Apache Airflow 是一个开源平台,用于开发、调度和监控面向批处理的工作流程。Airflow 所有的工作流程都可以使用 Python 代码定义。Web 界面可以管理工作流程的状态。
前提条件
已配置登录账号在各个所需监控项目的角色(项目管理员|实例管理员|数据读写),账号权限详情,请参见 成员管理。
已安装 Apache Airflow,详细说明,请参考 Apache Airflow 官网。
获取 OB Cloud 云数据库连接串
在实例列表页面,展开目标事务型实例的信息,在目标租户下,单击 连接 > 获取连接串。
在弹出框中选择 使用公共网络。
在 使用公网 IP 连接数据库 页面完成如下设置,生成连接串:
配置项 说明 添加 IP 地址 单击添加,将您的出口 IP 添加至白名单。 下载证书 (可选)单击下载证书,下载 CA 证书并完成认证。 连接租户 数据库 单击下拉框后,单击 + 创建数据库,根据提示完成数据库创建。 账号 单击下拉框后,单击 + 创建账号,根据提示完成账号创建。 连接方式 选择 MySQL CLI 作为连接方式。 注意
创建账号后,请您妥善记录创建账号时生成的密码。
在 Airflow 中添加 OB Cloud 数据源进行连接
打开 Airflow Web UI。
导航到 Admin -> Connections。
点击 "+" 号来添加一个新的连接。
填写以下字段:
配置项 说明 Connection Id ob(可以是任何标识符)。 Connection Type MySQL Host 取自连接串中 -h参数,OB Cloud 云数据库连接地址,例如,t5******.aws-ap-southeast-1.oceanbase.cloud。Schema 取自连接串中 -D参数,需要访问的数据库名称。Login 取自连接串中 -u参数,账号名称,例如:test。Password 取自连接串中 -p参数,账号密码。Port 取自连接串中 -P参数,OB Cloud 云数据库连接端口。创建成功后,即可在 Airflow 任务中通过引用
ob这个 Connection Id,来访问 OceanBase 数据库。
Airflow 任务示例
在 Airflow 中添加 OceanBase 数据库之后,可以通过编写以下内容,实现 AirFlow 从 OceanBase 数据库中读取数据并打印。
在 Airflow 安装目录下的 dags 文件夹中新建 query.py 文件,并编辑如下内容。
from airflow import DAG from airflow.utils.dates import days_ago from airflow.providers.mysql.hooks.mysql import MySqlHook from airflow.operators.python import PythonOperator default_args = { 'owner': 'airflow', 'retries': 0, } def fetch_and_print_data(): hook = MySqlHook(mysql_conn_id='ob') sql = "SELECT * FROM person LIMIT 1;" connection = hook.get_conn() cursor = connection.cursor() cursor.execute(sql) rows = cursor.fetchall() for row in rows: print(row) with DAG( dag_id='sql_query', default_args=default_args, schedule_interval='@daily', start_date=days_ago(1), catchup=False, ) as dag: run_and_print = PythonOperator( task_id='run_and_print', python_callable=fetch_and_print_data, ) run_and_print执行
airflow tasks test sql_query run_and_print,即可打印 person 表中的第一条数据。