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 プロトコルを使用する場合、データは以下の処理が必要です:
- サーバー上で行ベースのバイト列にシリアライズされる。
- ネットワーク経由で転送される。
- ターゲット構造に逆シリアライズされる(多くの場合、列形式への再変換が必要)。
この3ステップのプロセスにより、以下の問題が生じます:
- 高いシリアライゼーション/デシリアライゼーションのオーバーヘッド。
- 複雑なデータ変換。
- データ量に比例して増大するレイテンシ。
Arrow Flight SQL との統合は、以下の方法でこれらの問題を解決します:
- StarRocks の実行エンジンからクライアントまで、エンドツーエンドで列指向フォーマットを維持する。
- 分析ワークロード向けに最適化された Apache Arrow のインメモリ列指向表現を活用する。
- Arrow Flight のプロトコルを高速転送に使用し、中間変換なしで効率的なストリーミングを実現する。

この設計により、真のゼロコピー転送が実現され、従来の方法よりも高速かつリソース効率に優れています。
さらに、StarRocks は Arrow Flight SQL 向けのユニバーサル JDBC ドライバーを提供しており、アプリケーションは JDBC 互換性や他の Arrow Flight 対応システムとの相互運用性を犠牲にすることなく、この高性能転送パスを採用できます。
BE ノードがクライアントから直接アクセスできないデプロイ環境(プライベートネットワークや Kubernetes クラスターなど)向けに、StarRocks は Arrow Flight プロキシ機能を提供しています。有効にすると、FE がプロキシとして機能し、BE ノードからクライアントへ Arrow データをルーティングすることで、ネットワークトポロジーの制約に対応しながら列指向転送のメリットを維持します。このプロキシモードはわずかなパフォーマンスオーバーヘッドが発生しますが、BE への直接接続が利用できない環境でも Arrow Flight SQL アクセスを可能にします。
パフォーマンス比較
包括的なテストにより、データ取得速度の大幅な改善が実証されています。さまざまなデータ型(整数、浮動小数点、文字列、ブール値、混合カラム)において、Arrow Flight SQL は従来の PyMySQL および Pandas の read_sql インターフェースを一貫して上回りました。主な結果は以下のとおりです:
- 1,000 万行の整数データの読み取りでは、実行時間が約 35 秒から 0.4 秒に短縮(約 85 倍高速化)。
- 混合カラムテーブルでは、パフォーマンス改善が 160 倍の高速化に達した。
- 比較的単純なクエリ(例:単一の文字列カラム)でも、パフォーマンス向上は 12 倍を超えた。
平均して、Arrow Flight SQL は以下を達成しました:
- クエリの複雑さとデータ型に応じて、20 倍から 160 倍の転送時間の高速化。
- 冗長なシリアライゼーションステップの排除により、CPU およびメモリ使用量が明確に削減。
これらのパフォーマンス向上は、より高速なダッシュボード、より応答性の高いデータサイエンスワークフロー、そしてリアルタイムでより大規模なデータセットを分析する能力として直接反映されます。
クライアントコードでこれらの数値に到達する方法の詳細な内訳(JDBCアクセサーメソッド、生のVectorSchemaRoot消費、Parquetライター)、およびMySQL JDBCに対する各チューニングステップの測定済みスピードアップについては、Arrow Flight SQL ベストプラクティス.
使用方 法
Arrow Flight SQLプロトコルを介してPython ADBAドライバーを使用して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設定項目
JAVA_OPTSをfe.confで指定します。この方法では、他の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 Proxyの設定(オプション)
BEノードがクライアントアプリケーションから直接アクセスできない場合(例えば、プライベートネットワークやKubernetes環境にデプロイされている場合)、FE上でArrow Flightプロキシ機能を有効にして、BEノードからのデータをFEを通じてルーティングできます。
プロキシ機能は2つのグローバル変数によって制御されます:
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のホストアドレス
- StarRocks FEでArrow Flightがリッスンに使用するポート
- 必要な権限を持つ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(
🟢 Result:
StatusResult
0
⏱️ Execution time: 0.025 seconds
🟡 SQL:
SHOW DATABASES;
🟢 Result:
Database
_statistics_
hits
information_schema
sys
⏱️ Execution time: 0.014 seconds
🟡 SQL:
CREATE DATABASE sr_arrow_flight_sql;
🟢 Result:
StatusResult
0
⏱️ Execution time: 0.012 seconds