Apache Airflow
通过使用 DAG(有向无环图)和 SQL 操作符来实现与 StarRocks 的数据工作流编排和调度。使用 Airflow 进行数据导入和转换时,可以使用 SQLExecuteQueryOperator 和 MySQLHook,无需任何实现或复杂配置。
Apache Airflow GitHub 仓库。
支持的功能
- 通过 MySQL 协议执行 SQL
- 连接管理
- 事务支持
- 参数化查询
- 任务依赖
- 重试逻辑
安装
前提条件
- Apache Airflow 2.0+ 或 3.0+
- Python 3.8+
- 访问 StarRocks 集群(请参阅 快速入门指南)
安装
要使用 StarRocks,需要安装 MySQL 提供程序包,因为 StarRocks 使用 MySQL 协议。
pip install apache-airflow-providers-mysql
通过检查已安装的提供程序来验证安装:
airflow providers list
输出中应列出 apache-airflow-providers-mysql。
配置
创建 StarRocks 连接
在 Airflow UI 或通过环境变量创建一个 StarRocks 连接。连接的名称将在 DAG 中使用。
通过 Airflow UI
- 导航到 Admin > Connections
- 点击 + 按钮添加新连接
- 配置连接:
- Connection Id:
starrocks_default - Connection Type: MySQL
- Host:
your-starrocks-host.com - Schema:
your_database - Login:
your_username - Password:
your_password - Port:
9030
通过 Airflow CLI
airflow connections add 'starrocks_default' \
--conn-type 'mysql' \
--conn-host 'your-starrocks-host.com' \
--conn-schema 'your_database' \
--conn-login 'your_username' \
--conn-password 'your_password' \
--conn-port 9030
使用示例
这些示例展示了将 StarRocks 与 Airflow 集成的常见模式。每个示例都基于核心概念,同时展示了不同的数据导入、转换和工作流编排方法。
您将学习到:
- 数据导入:高效地从 CSV 文件和云存储中导入数据到 StarRocks
- 数据转换:执行 SQL 查询并使用 Python 处理结果
- 高级模式:实现增量导入、异步操作和查询优化
- 生产最佳实践:优雅地处理错误并构建可靠的管道
所有示例都使用 快速入门指南 中描述的崩溃数据表。
数据导入
流式数据导入
使用 StarRocks Stream Load API 高效导入大型 CSV 文件。Stream Load 是推荐的方法:
- 高吞吐量数据导入(支持并行导入)
- 导入数据时进行列转换和过滤
对于大型数据集,Stream Load 比 INSERT INTO VALUES 语句提供更好的性能,并包含内置的错误容忍功能。注意,这需要 CSV 文件可以在 Airflow 工作器的文件系统上访问。
from airflow.sdk import dag, task
from airflow.hooks.base import BaseHook
from datetime import datetime
import requests
from requests.auth import HTTPBasicAuth
from urllib.parse import urlparse
class PreserveAuthSession(requests.Session):
"""
自定义会话,保留重定向时的授权头。
StarRocks FE 可能会将 Stream Load 请求重定向到 BE 节点。
"""
def rebuild_auth(self, prepared_request, response):
old = urlparse(response.request.url)
new = urlparse(prepared_request.url)
# 仅在重定向到相同主机名时保留授权
if old.hostname == new.hostname:
prepared_request.headers["Authorization"] = response.request.headers.get("Authorization")
@dag(
dag_id="starrocks_stream_load_example",
schedule=None,
start_date=datetime(2024, 1, 1),
catchup=False,
tags=["starrocks", "stream_load", "example"],
)
def starrocks_stream_load_example():
@task
def load_csv_to_starrocks():
# 配置
DATABASE = "quickstart"
TABLE = "crashdata"
CSV_PATH = "/path/to/crashdata.csv"
conn = BaseHook.get_connection("starrocks_default")
url = f"http://{conn.host}:{conn.port}/api/{DATABASE}/{TABLE}/_stream_load"
# 生成唯一标签
from airflow.sdk import get_current_context
context = get_current_context()
execution_date = context['logical_date'].strftime('%Y%m%d_%H%M%S')
label = f"{TABLE}_load_{execution_date}"
headers = {
"label": label,
"column_separator": ",",
"skip_header": "1",
"max_filter_ratio": "0.1", # 允许最多 10% 的错误率
"Expect": "100-continue",
"columns": """
tmp_CRASH_DATE, tmp_CRASH_TIME,
CRASH_DATE=str_to_date(concat_ws(' ', tmp_CRASH_DATE, tmp_CRASH_TIME), '%m/%d/%Y %H:%i'),
BOROUGH, ZIP_CODE, LATITUDE, LONGITUDE, LOCATION,
ON_STREET_NAME, CROSS_STREET_NAME, OFF_STREET_NAME,
NUMBER_OF_PERSONS_INJURED, NUMBER_OF_PERSONS_KILLED,
NUMBER_OF_PEDESTRIANS_INJURED, NUMBER_OF_PEDESTRIANS_KILLED,
NUMBER_OF_CYCLIST_INJURED, NUMBER_OF_CYCLIST_KILLED,
NUMBER_OF_MOTORIST_INJURED, NUMBER_OF_MOTORIST_KILLED,
CONTRIBUTING_FACTOR_VEHICLE_1, CONTRIBUTING_FACTOR_VEHICLE_2,
CONTRIBUTING_FACTOR_VEHICLE_3, CONTRIBUTING_FACTOR_VEHICLE_4,
CONTRIBUTING_FACTOR_VEHICLE_5, COLLISION_ID,
VEHICLE_TYPE_CODE_1, VEHICLE_TYPE_CODE_2,
VEHICLE_TYPE_CODE_3, VEHICLE_TYPE_CODE_4, VEHICLE_TYPE_CODE_5
""".replace("\n", "").replace(" ", ""),
}
session = PreserveAuthSession()
with open(CSV_PATH, "rb") as f:
response = session.put(
url,
headers=headers,
data=f,
auth=HTTPBasicAuth(conn.login, conn.password or ""),
timeout=3600,
)
result = response.json()
print(f"\nStream Load Response:")
print(f" Status: {result.get('Status')}")
print(f" Loaded Rows: {result.get('NumberLoadedRows', 0):,}")
if result.get("Status") == "Success":
return result
else:
error_msg = result.get("Message", "Unknown error")
raise Exception(f"Stream Load failed: {error_msg}")
load_csv_to_starrocks()
starrocks_stream_load_example()