通过 Arrow Flight SQL 与 StarRocks 交互
从 v3.5.1 起,StarRocks 支持通过 Apache Arrow Flight SQL 协议进行连 接。
概述
借助 Arrow Flight SQL 协议,您可以执行普通的 DDL、DML、DQL 语句,并使用 Python 代码或 Java 代码通过 Arrow Flight SQL ADBC 或 JDBC 驱动读取大规模数据。
该方案从 StarRocks 列式执行引擎到客户端建立了一条全列式数据传输管道,消除了传统 JDBC 和 ODBC 接口中常见的频繁行列转换和序列化开销,使 StarRocks 能够以零拷贝、低延迟、高吞吐量的方式传输数据。
应用场景
Arrow Flight SQL 集成使 StarRocks 特别适用于:
- 数据科学工作流,其中 Pandas 和 Apache Arrow 等工具需要列式数据。
- 数据湖分析,需要对海量数据集进行高吞吐量、低延迟的访问。
- 机器学习,快速迭代和处理速度至关重要。
- 实时分析平台,必须以最小延迟交付数据。
借助 Arrow Flight SQL,您可以获益于:
- 端到端列式数据传输,消除列式与行式格式之间的高代价转换。
- 零拷贝数据移动,降低 CPU 和内存开销。
- 低延迟和极高吞吐量,加速分析和响应速度。
技术方案
传统上,StarRocks 在内部以列式 Block 结构组织查询结果。然而,在使用 JDBC、ODBC 或 MySQL 协议时,数据必须经过:
- 在服务端序列化为基于行的字节流。
- 通过网络传输。
- 反序列化还原为目标结构(通常需要重新转换为列式格式)。
这三步过程会带来:
- 高序列化/反序列化开销。
- 复杂的数据转换。
- 随数据量增长而增加的延迟。
与 Arrow Flight SQL 的集成通过以下方式解决了这些问题:
- 从 StarRocks 执行引擎直接到客户端,全程保留列式格式。
- 利用 Apache Arrow 的内存列式表示,该表示针对分析工作负载进行了优化。
- 使用 Arrow Flight 协议进行高速传输,无需中间转换即可实现高效流式传输。

该设计提供了真正的零拷贝传输,比传统方法更快且更节省资源。
此外,StarRocks 为 Arrow Flight SQL 提供了通用 JDBC 驱动,使应用程序无需牺牲 JDBC 兼容性或与其他支持 Arrow Flight 的系统的互操作性,即可采用这条高性能传输路径。
对于 BE 节点无法被客户端直接访问的部署场景(例如私有网络或 Kubernetes 集群),StarRocks 提供了 Arrow Flight 代理功能。启用后,FE 可作为代理将 Arrow 数据从 BE 节点路由到客户端,在满足网络拓扑约束的同时保留列式传输的优势。该代理模式会带来少量性能开销,但可在无法直接连接 BE 的环境中启用 Arrow Flight SQL 访问。
性能对比
综合测试表明,数据检索速度有了显著提升。 在各种数据类型(整数、浮点数、字符串、布尔值及混合列)下,Arrow Flight SQL 的性能始终优于传统的 PyMySQL 和 Pandas read_sql 接口。主要结果包括:
- 读取 1000 万行整数数据时,执行时间从约 35 秒降至 0.4 秒(约快 85 倍)。
- 对于混合列表,性能提升达到 160 倍加速。
- 即使在较简单的查询中(例如单字符串列),性能提升也超过 12 倍。
平均而言,Arrow Flight SQL 实现了:
- 20 倍至 160 倍的传输速度提升,具体取决于查询复杂度和数据类型。
- 由于消除了冗余序列化步骤,CPU 和内存使用量明显降低。
这些性能提升直接转化为更快的仪表盘响应、更流畅的数据科学工作流,以及实时分析更大规模数据集的能力。
有关如何在客户端代码中实现这些数字的详细分解——JDBC 访问器方法、原始 VectorSchemaRoot 消耗、Parquet 写入器——以及每个调优步骤相对于 MySQL JDBC 的实测加速效果,请参阅 Arrow Flight SQL 最佳实践.
使用方法
按照以下步骤,通过 Arrow Flight SQL 协议使用 Python ADBC Driver 连接 StarRocks 并与之交互。完整代码示例请参阅 附录。
前提条件:Python 3.9 或更高版本。
步骤 1:安装库
使用 pip 从 PyPI 安装 adbc_driver_manager 和 adbc_driver_flightsql:
pip install adbc_driver_manager
pip install adbc_driver_flightsql
将以下模块或库导入您的代码:
- 必需库:
import adbc_driver_manager
import adbc_driver_flightsql.dbapi as flight_sql
- 可选模块(用于提升可用性和调试体验):
import pandas as pd # Optional: for better result display using DataFrame
import traceback # Optional: for detailed error traceback during SQL execution
import time # Optional: for measuring SQL execution time
步骤 2:连接到 StarRocks
-
如果您想通过命令行启动 FE 服务,可以使用以下任意一种方式:
-
指定环境变量
JAVA_TOOL_OPTIONS。export JAVA_TOOL_OPTIONS="--add-opens=java.base/java.nio=org.apache.arrow.memory.core,ALL-UNNAMED" -
在 fe.conf 中指定 FE 配置项
JAVA_OPTS,这样可以追加其他JAVA_OPTS值。JAVA_OPTS="--add-opens=java.base/java.nio=org.apache.arrow.memory.core,ALL-UNNAMED ..."
-
-
如果您想在 IntelliJ IDEA 中运行服务,必须在
Run/Debug Configurations的Build and run中添加以下选项:--add-opens=java.base/java.nio=org.apache.arrow.memory.core,ALL-UNNAMED
配置 StarRocks
在通过 Arrow Flight SQL 连接 StarRocks 之前,您必须先配置 FE 和 BE 节点,以确保 Arrow Flight SQL 服务已启用并监听指定端口。
在 FE 配置文件 fe.conf 和 BE 配置文件 be.conf 中,将 arrow_flight_port 设置为可用端口。修改配置文件后,重启 FE 和 BE 服务以使修改生效。
FE 和 BE 必须设置不同的 arrow_flight_port。
示例:
// fe.conf
arrow_flight_port = 9408
// be.conf
arrow_flight_port = 9419
配置 Arrow Flight 代理(可选)
如果您的 BE 节点无法从客户端应用程序直接访问(例如,部署在私有网络或 Kubernetes 环境中),您可以在 FE 上启用 Arrow Flight 代理功能,将 BE 节点的数据通过 FE 进行路由。
代理功能由两个全局变量控制:
arrow_flight_proxy_enabled:控制是否启用代理模式。默认值为true。启用后会有轻微的性能开销。arrow_flight_proxy:指定代理主机名。若为空(默认值),则当前 FE 节点充当代理。如果使用不同的代理端点,可将其设置为特定主机名。
要为所有会话全局配置这些变量:
-- 启用或禁用代理模式(默认启用)
SET GLOBAL arrow_flight_proxy_enabled = true;
-- 设置特定的代理主机名(可选,默认为当前 FE)
SET GLOBAL arrow_flight_proxy = 'your-proxy-hostname:Port';
- 代理功能默认启用,与直接连接 BE 相比,可能导致吞吐量降低 8-10%。如果您的客户端可以直接访问 BE 节点的网络,或者 FE 侧内存资源有限,可以禁用代理以获得最佳性能:
SET GLOBAL arrow_flight_proxy_enabled = false;。 - 当
arrow_flight_proxy为空时,票据将自动通过客户端最初连接的 FE 节点进行路由。 - 重要:
arrow_flight_proxy和arrow_flight_proxy_enabled设置应使用SET GLOBAL进行全局配置。不支持会话级别的设置。 - 需要重启会话:更改代理设置仅影响新会话。现有的 Arrow Flight SQL 会话将继续使用其原始设置,直到重新连接。
建立连接
在客户端,使用以下信息创建 Arrow Flight SQL 客户端:
- StarRocks FE 的主机地址
- Arrow Flight 在 StarRocks FE 上监听所使用的端口
- 具有必要权限的 StarRocks 用户的用户名和密码
示例:
FE_HOST = "127.0.0.1"
FE_PORT = 9408
conn = flight_sql.connect(
uri=f"grpc://{FE_HOST}:{FE_PORT}",
db_kwargs={
adbc_driver_manager.DatabaseOptions.USERNAME.value: "root",
adbc_driver_manager.DatabaseOptions.PASSWORD.value: "",
}
)
cursor = conn.cursor()
连接建立后,您可以通过返回的 Cursor 执行 SQL 语句与 StarRocks 进行交互。
步骤 3.(可选)预定义工具函数
这些函数用于格式化输出、统一格式并简化调试。您可以在代码中选择性地定义它们以进行测试。
# =============================================================================
# 用于更好地格式化输出和执行 SQL 的工具函数
# =============================================================================
# 打印章节标题
def print_header(title: str):
"""
Print a section header for better readability.
"""
print("\n" + "=" * 80)
print(f"🟢 {title}")
print("=" * 80)
# 打印正在执行的 SQL 语句
def print_sql(sql: str):
"""
Print the SQL statement before execution.
"""
print(f"\n🟡 SQL:\n{sql.strip()}")
# 打印结果 DataFrame
def print_result(df: pd.DataFrame):
"""
Print the result DataFrame in a readable format.
"""
if df.empty:
print("\n🟢 Result: (no rows returned)\n")
else:
print("\n🟢 Result:\n")
print(df.to_string(index=False))
# 打印错误堆栈信息
def print_error(e: Exception):
"""
Print the error traceback if SQL execution fails.
"""
print("\n🔴 Error occurred:")
traceback.print_exc()
# 执行 SQL 语句并打印结果
def execute(sql: str):
"""
Execute a SQL statement and print the result and execution time.
"""
print_sql(sql)
try:
start = time.time() # Optional: start time for execution time measurement
cursor.execute(sql)
result = cursor.fetchallarrow() # Arrow Table
df = result.to_pandas() # Optional: convert to DataFrame for better display
print_result(df)
print(f"\n⏱️ Execution time: {time.time() - start:.3f} seconds")
except Exception as e:
print_error(e)
步骤 4. 与 StarRocks 交互
本节将引导您完成一些基本操作,例如创建表、加载数据、检查表元数据、设置变量以及运行查询。
以下列出的输出示例基于前述步骤中描述的可选模块和工具函数实现。
-
创建数据库和表以加载数据,并检查表结构。
# 第 1 步:删除并创建数据库
print_header("Step 1: Drop and Create Database")
execute("DROP DATABASE IF EXISTS sr_arrow_flight_sql FORCE;")
execute("SHOW DATABASES;")
execute("CREATE DATABASE sr_arrow_flight_sql;")
execute("SHOW DATABASES;")
execute("USE sr_arrow_flight_sql;")
# 第 2 步:创建表
print_header("Step 2: Create Table")
execute("""
CREATE TABLE sr_arrow_flight_sql_test
(
k0 INT,
k1 DOUBLE,
k2 VARCHAR(32) NULL DEFAULT "" COMMENT "",
k3 DECIMAL(27,9) DEFAULT "0",
k4 BIGINT NULL DEFAULT '10',
k5 DATE
)
DISTRIBUTED BY HASH(k5) BUCKETS 5
PROPERTIES("replication_num" = "1");
""")
execute("SHOW CREATE TABLE sr_arrow_flight_sql_test;")示例输出:
================================================================================
🟢 Step 1: Drop and Create Database
================================================================================
🟡 SQL:
DROP DATABASE IF EXISTS sr_arrow_flight_sql FORCE;
/Users/starrocks/test/venv/lib/python3.9/site-packages/adbc_driver_manager/dbapi.py:307: Warning: Cannot disable autocommit; conn will not be DB-API 2.0 compliant
warnings.warn(