Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

DataStoreの実行モデル

DataStoreの遅延評価モデルを理解することは、DataStoreを効果的に活用し、最適なパフォーマンスを引き出すうえで重要です。

遅延評価

DataStore は 遅延評価 を採用しています。つまり、操作はすぐには実行されず、いったん記録された後、最適化された SQL クエリへコンパイルされます。実行されるのは、結果が実際に必要になったときだけです。

例: 遅延評価と即時評価

from pathlib import Path
Path("sales.csv").write_text("""\
region,product,category,amount,quantity,price,date,order_id
East,Widget,Electronics,5200,10,120,2024-01-15,1001
West,Gadget,Electronics,800,5,160,2024-02-20,1002
East,Gizmo,Home,6500,3,100,2024-03-10,1003
North,Widget,Electronics,4500,6,150,2024-06-18,1004
West,Gadget,Electronics,2000,8,250,2024-09-14,1005
""")

from chdb import datastore as pd

ds = pd.read_csv("sales.csv")

# These operations are NOT executed yet
result = (ds
    .filter(ds['amount'] > 1000)    # Recorded, not executed
    .select('region', 'amount')      # Recorded, not executed
    .groupby('region')               # Recorded, not executed
    .agg({'amount': 'sum'})          # Recorded, not executed
    .sort('sum', ascending=False)    # Recorded, not executed
)

# Still no execution - just building the query plan
print(result.to_sql())
# SELECT region, SUM(amount) AS sum
# FROM file('sales.csv', 'CSVWithNames')
# WHERE amount > 1000
# GROUP BY region
# ORDER BY sum DESC

# NOW execution happens
df = result.to_df()  # <-- Triggers execution

遅延評価の利点

  1. クエリ最適化: 複数の操作が、最適化された単一のSQLクエリにまとめてコンパイルされます
  2. フィルタのプッシュダウン: フィルタはデータソースレベルで適用されます
  3. カラムプルーニング: 必要なカラムだけが読み込まれます
  4. 実行時の判断: 実行エンジンは実行時に選択できます
  5. プランの確認: 実行前にクエリを確認したりデバッグしたりできます

実行のトリガー

実際の値が必要になると、自動的に実行がトリガーされます:

自動トリガー

トリガー 説明
print() / repr() print(ds) 結果を表示
len() len(ds) 行数を取得
.columns ds.columns カラム名を取得
.dtypes ds.dtypes カラム型を取得
.shape ds.shape 形状を取得
.index ds.index 行インデックスを取得
.values ds.values NumPy配列を取得
Iteration for row in ds 行を順に処理
to_df() ds.to_df() pandasに変換
to_pandas() ds.to_pandas() to_df のエイリアス
to_dict() ds.to_dict() dictに変換
to_numpy() ds.to_numpy() 配列に変換
.equals() ds.equals(other) DataStoresを比較

例:

# All these trigger execution
print(ds)              # Display
len(ds)                # 1000
ds.columns             # Index(['name', 'age', 'city'])
ds.shape               # (1000, 3)
list(ds)               # List of values
ds.to_df()             # pandas DataFrame

遅延実行のまま維持される操作

Operation Returns Description
filter() DataStore WHERE 句を追加
select() DataStore 選択するカラムを追加
sort() DataStore ORDER BY を追加
groupby() LazyGroupBy GROUP BY を設定
join() DataStore JOIN を追加
ds['col'] ColumnExpr カラム参照
ds[['col1', 'col2']] DataStore カラムの選択

例:

# These do NOT trigger execution - they stay lazy
result = ds.filter(ds['age'] > 25)      # Returns DataStore
result = ds.select('name', 'age')        # Returns DataStore
result = ds['name']                      # Returns ColumnExpr
result = ds.groupby('city')              # Returns LazyGroupBy

3 段階の実行

DataStore の操作は、3 段階の実行モデルに沿って進行します。

フェーズ 1: SQLクエリの構築 (遅延実行)

SQLで表現できる操作が蓄積されます:

result = (ds
    .filter(ds['status'] == 'active')   # WHERE
    .select('user_id', 'amount')         # SELECT
    .groupby('user_id')                  # GROUP BY
    .agg({'amount': 'sum'})              # SUM()
    .sort('sum', ascending=False)        # ORDER BY
    .limit(10)                           # LIMIT
)
# All compiled into one SQL query

フェーズ2: 実行ポイント

トリガーが発生すると、それまでに蓄積されたSQLが実行されます。

# Execution triggered here
df = result.to_df()  
# The single optimized SQL query runs now

フェーズ 3: DataFrame の操作 (ある場合)

実行後に pandas 固有の操作を続けて適用する場合:

# Mixed operations
result = (ds
    .filter(ds['amount'] > 100)          # Phase 1: SQL
    .to_df()                             # Phase 2: Execute
    .pivot_table(...)                    # Phase 3: pandas
)

実行計画を表示する

explain() を使用すると、実行される内容を確認できます。

Querypython
ds = pd.read_csv("sales.csv")

query = (ds
    .filter(ds['amount'] > 1000)
    .groupby('region')
    .agg({'amount': ['sum', 'mean']})
)

# View execution plan
query.explain()
Responsetext
Pipeline:
  1. Source: file('sales.csv', 'CSVWithNames')
  2. Filter: amount > 1000
  3. GroupBy: region
  4. Aggregate: sum(amount), avg(amount)

Generated SQL:
SELECT region, SUM(amount) AS sum, AVG(amount) AS mean
FROM file('sales.csv', 'CSVWithNames')
WHERE amount > 1000
GROUP BY region

詳細情報を表示するには verbose=True を使用します:

query.explain(verbose=True)

詳しいドキュメントについては、Debugging: explain()を参照してください。


キャッシュ

DataStore は、重複するクエリの実行を避けるために、実行結果をキャッシュします。

キャッシュの仕組み

from pathlib import Path
Path("data.csv").write_text("""\
name,age,city,salary,department
Alice,25,NYC,55000,Engineering
Bob,30,LA,65000,Product
Charlie,35,NYC,80000,Engineering
Diana,28,SF,70000,Design
Eve,42,NYC,95000,Product
""")

ds = pd.read_csv("data.csv")
result = ds.filter(ds['age'] > 25)

# First access - executes query
print(result.shape)  # Executes and caches

# Second access - uses cache
print(result.columns)  # Uses cached result

# Third access - uses cache
df = result.to_df()  # Uses cached result

キャッシュの無効化

DataStore に変更を加える操作が行われると、キャッシュは無効化されます:

result = ds.filter(ds['age'] > 25)
print(result.shape)  # Executes, caches

# New operation invalidates cache
result2 = result.filter(result['city'] == 'NYC')
print(result2.shape)  # Re-executes (different query)

キャッシュの手動制御

# Clear cache
ds.clear_cache()

# Disable caching
from chdb.datastore.config import config
config.set_cache_enabled(False)

SQL と Pandas の操作を組み合わせる

DataStore は、SQL と pandas を組み合わせた操作を適切に処理します。

SQL互換の操作

以下はSQLに変換されます:

  • filter(), where()
  • select()
  • groupby(), agg()
  • sort(), orderby()
  • limit(), offset()
  • join(), union()
  • distinct()
  • カラム操作 (算術演算、比較、文字列メソッド)

Pandas のみの操作

以下は実行がトリガーされ、pandas が使用されます。

  • カスタム関数を使った apply()
  • 複雑な集計を行う pivot_table()
  • stack()unstack()
  • 実行済みの DataFrame に対する操作

ハイブリッドパイプライン

# SQL phase
result = (ds
    .filter(ds['amount'] > 100)      # SQL
    .groupby('category')              # SQL
    .agg({'amount': 'sum'})           # SQL
)

# Execution + pandas phase
result = (result
    .to_df()                          # Execute SQL
    .pivot_table(...)                 # pandas operation
)

実行エンジンの選択

DataStore では、異なるエンジンを使って操作を実行できます。

自動モード (既定)

from chdb.datastore.config import config

config.set_execution_engine('auto')  # Default
# Automatically selects best engine per operation

chDB エンジンを強制的に使用する

config.set_execution_engine('chdb')
# All operations use ClickHouse SQL

pandas Engineの使用を強制する

config.set_execution_engine('pandas')
# All operations use pandas

詳しくは、設定: 実行エンジンを参照してください。


パフォーマンスへの影響

良い例: 早い段階でフィルタする

# Good: Filter in SQL, then aggregate
result = (ds
    .filter(ds['date'] >= '2024-01-01')  # Reduces data early
    .groupby('category')
    .agg({'amount': 'sum'})
)

悪い例: フィルタを後で適用する

# Bad: Aggregate all, then filter
result = (ds
    .groupby('category')
    .agg({'amount': 'sum'})
    .to_df()
    .query('sum > 1000')  # Pandas filter after aggregation
)

良い例: 早い段階でカラムを絞り込む

# Good: Select columns in SQL
result = (ds
    .select('user_id', 'amount', 'date')
    .filter(ds['date'] >= '2024-01-01')
    .groupby('user_id')
    .agg({'amount': 'sum'})
)

良い例: SQLに任せる

# Good: Complex aggregation in SQL
result = (ds
    .groupby('category')
    .agg({
        'amount': ['sum', 'mean', 'count'],
        'quantity': 'sum'
    })
    .sort('sum', ascending=False)
    .limit(10)
)
# One SQL query does everything

# Bad: Multiple separate queries
sums = ds.groupby('category')['amount'].sum().to_df()
means = ds.groupby('category')['amount'].mean().to_df()
# Two queries instead of one

ベストプラクティスのまとめ

  1. 実行前に一連の操作を組み立てる - クエリ全体を構築してから、一度だけトリガーします
  2. 早い段階で絞り込む - データの発生元でデータ量を減らします
  3. 必要なカラムだけを選択する - カラムプルーニングによりパフォーマンスが向上します
  4. 実行内容を理解するために explain() を使う - 実行前にデバッグします
  5. 集計は SQL に任せる - ClickHouse はこの処理に最適化されています
  6. 実行のトリガーを意識する - 意図しない早期実行を避けます
  7. キャッシュを適切に使う - cache が無効になるタイミングを理解します
Navigation