•20 min read

2026年のカラム型データ: Apache Arrowメモリレイアウト、Parquetストレージ、DuckDBベクタライゼーション

2026年のカラム型データ: Apache Arrowメモリレイアウト、Parquetストレージ、DuckDBベクタライゼーション

カラム型データフォーマットは、最新の分析システムの基盤であり、ペタバイト規模のデータセットに対して秒以下のクエリパフォーマンスを可能にします。このガイドでは、永続ストレージ用のApache Parquet、インメモリ処理用のApache Arrow、およびDuckDBのベクトル化実行エンジンの相互作用を分析し、これらのテクノロジーがどのように連携して高性能な分析データスタックを形成するかを説明します。

Audio Briefing
0:00 / 0:00

Apache Parquet: ディスク最適化されたカラム型ストレージ

Apache Parquetは、ビッグデータエコシステムにおけるカラム型データストレージの事実上の標準です。その設計は、分析ワークロードにとって不可欠な効率的なディスクI/Oと述語プッシュダウンを優先しています。

Parquetファイル構造

Parquetファイルは、データの水平パーティションである行グループで構成されています。各行グループ内では、データはカラムごとに保存されます。この構造により、クエリに必要なカラムのみを読み取ることができ、I/Oを大幅に削減します。

Parquet File
├── File Metadata (footer)
│   └── Schema, Row Group Metadata, Column Metadata
├── Row Group 1
│   ├── Column 1 Chunk
│   │   ├── Page 1 (Data, Dictionary, Index)
│   │   ├── Page 2
│   │   └── ...
│   ├── Column 2 Chunk
│   └── ...
├── Row Group 2
│   ├── Column 1 Chunk
│   └── ...
└── ...

エンコーディングと圧縮

Parquetは、ストレージフットプリントを最小限に抑え、読み取りパフォーマンスを向上させるために、さまざまなエンコーディングスキームと圧縮アルゴリズムを採用しています。

  • 辞書エンコーディング(Dictionary Encoding): カーディナリティの低いカラムの場合、一意の値が辞書に保存され、カラムデータにはこの辞書へのインデックスが保存されます。これは、文字列データやカテゴリデータに非常に効果的です。
  • ランレングスエンコーディング(Run Length Encoding (RLE)): 繰り返し値のシーケンスに効率的です。
  • ビットパックエンコーディング(Bit-Packed Encoding): 整数値の場合、必要な最小ビット数に値をパックします。
  • 圧縮: Snappy、Gzip、Zstd、Brotli。Zstdは、圧縮率と速度の優れたバランスを提供します。

述語プッシュダウンとカラム統計

Parquetのファイル形式には、行グループ内の各カラムチャンクの統計(最小/最大値、NULLカウント)が含まれています。クエリエンジンはこれらの統計を利用して、クエリの述語を満たすことができない行グループやページ全体をプルーニングし、不要なI/Oを回避します。

クエリSELECT sum(value) FROM table WHERE timestamp > '2023-01-01'を考えてみましょう。行グループのtimestampカラム統計がその最大タイムスタンプが2022-12-31であることを示している場合、その行グループ全体がスキップされます。

import pandas as pd
import pyarrow as pa
import pyarrow.parquet as pq
import numpy as np
import os

# Generate sample data
data = {
    'id': np.arange(100000),
    'timestamp': pd.to_datetime(pd.date_range('2022-01-01', periods=100000, freq='H')),
    'category': np.random.choice(['A', 'B', 'C', 'D'], 100000),
    'value': np.random.rand(100000) * 100
}
df = pd.DataFrame(data)

# Convert to PyArrow Table
table = pa.Table.from_pandas(df)

# Define Parquet file path
parquet_file = 'sample_data.parquet'

# Write to Parquet with specific options
# row_group_size: Controls the number of rows in each row group.
#                 Smaller row groups can improve pruning but increase metadata overhead.
#                 Larger row groups reduce metadata but might read more data than needed.
# compression: Zstd offers good balance.
# write_statistics: Essential for predicate pushdown.
pq.write_table(
    table,
    parquet_file,
    row_group_size=10000, # 10,000 rows per row group
    compression='zstd',
    write_statistics=True,
    data_page_size=8 * 1024 * 1024 # 8MB data page size
)

print(f"Parquet file '{parquet_file}' created.")

# Inspect Parquet file metadata
parquet_schema = pq.read_schema(parquet_file)
print("\nParquet Schema:")
print(parquet_schema)

parquet_metadata = pq.read_metadata(parquet_file)
print("\nParquet Metadata:")
print(parquet_metadata)

# Example of reading with predicate pushdown
# This query will only read row groups where 'timestamp' is after '2023-01-01'
# and only the 'timestamp' and 'value' columns.
start_time = pd.Timestamp('2023-01-01')
filtered_table = pq.read_table(
    parquet_file,
    columns=['timestamp', 'value'],
    filters=[('timestamp', '>', start_time)]
)

print(f"\nFiltered table (first 5 rows) after predicate pushdown and column projection:")
print(filtered_table.to_pandas().head())

# Clean up
os.remove(parquet_file)
Advertisement

Apache Arrow: インメモリ・カラム型フォーマット

Apache Arrowは、カラム型データのための言語に依存しない連続したメモリレイアウトを定義します。その主な目標は、データが異なるシステムやプロセス間を移動する際のシリアライズ/デシリアライズのオーバーヘッドを排除することです。

Arrowメモリレイアウト

Arrowの強みは、標準化されたフラットで連続したメモリバッファにあります。特定のカラムの場合、同じ型のすべての値が単一のバッファにまとめて保存されます。このレイアウトはキャッシュに優しく、効率的なSIMD(Single Instruction, Multiple Data)操作を可能にします。

単純な整数カラムの場合: [value1, value2, value3, ..., valueN]

NULL許容の整数カラムの場合: [validity_bitmap_buffer], [value1, value2, value3, ..., valueN] 有効性ビットマップは、どの値がNULLであるかを示します。

可変長型(例:文字列)の場合: [offset_buffer], [data_buffer], [validity_bitmap_buffer] オフセットバッファは、データバッファ内の各文字列の開始位置と長さを格納します。

import pyarrow as pa
import numpy as np

# Example: Integer Array
int_array = pa.array([1, 2, None, 4, 5], type=pa.int32())
print("Integer Array:")
print(int_array)
print(f"  Type: {int_array.type}")
print(f"  Buffers: {[b.to_pybytes() for b in int_array.buffers()]}") # Validity bitmap, Data buffer
print(f"  Is Null: {int_array.is_null()}")
print(f"  Values: {int_array.to_numpy()}")

# Example: String Array
string_array = pa.array(["apple", "banana", None, "cherry"], type=pa.string())
print("\nString Array:")
print(string_array)
print(f"  Type: {string_array.type}")
print(f"  Buffers: {[b.to_pybytes() for b in string_array.buffers()]}") # Validity bitmap, Offset buffer, Data buffer
print(f"  Is Null: {string_array.is_null()}")
print(f"  Values: {string_array.to_numpy()}")

# Example: RecordBatch (collection of arrays/columns)
schema = pa.schema([
    pa.field('id', pa.int64()),
    pa.field('name', pa.string()),
    pa.field('score', pa.float32())
])

record_batch = pa.RecordBatch.from_arrays(
    [
        pa.array([1, 2, 3, 4]),
        pa.array(["Alice", "Bob", "Charlie", "David"]),
        pa.array([90.5, 88.0, 92.1, 78.9])
    ],
    schema=schema
)

print("\nRecordBatch:")
print(record_batch)
print(f"  Number of columns: {record_batch.num_columns}")
print(f"  Number of rows: {record_batch.num_rows}")
print(f"  Column 'name' data buffer (first 10 bytes): {record_batch.column('name').buffers()[2].to_pybytes()[:10]}")

ゼロコピー読み取り

Arrowの最も重要な利点は、そのゼロコピー機能です。データがすでにArrow形式である場合、シリアライズやデシリアライズなしで、プロセス間、言語間、さらには異なるマシン間(Arrow Flight経由)で渡すことができます。これにより、CPUサイクルとレイテンシが大幅に削減されます。

例えば、ParquetファイルはArrow TableまたはRecordBatchに直接読み込むことができます。Parquetリーダーは、ディスク上のカラム型データをArrowのインメモリレイアウトにマッピングし、多くの場合、最小限のデータ変換で処理します。

DuckDB: ベクトル化クエリエンジン

DuckDBは、分析ワークロード向けに設計されたインプロセスSQL OLAPデータベースです。その核となる強みは、カラム型データ(多くの場合Arrow形式)で直接動作し、SIMD命令を活用して極限のパフォーマンスを実現するベクトル化クエリ実行エンジンにあります。

ベクトル化実行

DuckDBは、一度に1行ずつ(タプルごとの処理)ではなく、通常1024〜4096行のバッチ(ベクトル)でデータを処理します。各操作(例:フィルタ、結合、集計)はベクトル全体で動作するため、次のことが可能になります。

  1. 関数呼び出しのオーバーヘッドの削減: 単一の関数呼び出しで多くの値を処理します。
  2. CPUキャッシュ効率: データが順次アクセスされるため、キャッシュヒット率が高くなります。
  3. SIMD命令の活用: 最新のCPUは、SIMD命令(例:AVX2、AVX-512)を使用して、複数のデータポイントに対して同じ操作を同時に実行できます。DuckDBのカラム型レイアウトは、これに完全に適しています。

ParquetとArrowとの統合

DuckDBは、明示的なCREATE TABLEステートメントやデータ取り込みを必要とせずに、Parquetファイルを直接クエリできます。Parquetファイルを読み取り、それを内部のカラム型表現(Arrow互換)に変換し、そのベクトル化エンジンを使用してクエリを実行します。これにより、多くの分析タスクでETLステップが不要になります。

import duckdb
import pyarrow as pa
import pyarrow.parquet as pq
import pandas as pd
import numpy as np
import os
import time

# 1. Create a large Parquet file for demonstration
num_rows = 10_000_000
data = {
    'id': np.arange(num_rows),
    'timestamp': pd.to_datetime(pd.date_range('2020-01-01', periods=num_rows, freq='10s')),
    'category': np.random.choice(['A', 'B', 'C', 'D', 'E'], num_rows),
    'value': np.random.rand(num_rows) * 1000,
    'is_active': np.random.choice([True, False], num_rows, p=[0.7, 0.3])
}
df_large = pd.DataFrame(data)
table_large = pa.Table.from_pandas(df_large)
large_parquet_file = 'large_sample_data.parquet'
pq.write_table(table_large, large_parquet_file, row_group_size=100000, compression='zstd', write_statistics=True)
print(f"Created large Parquet file: {large_parquet_file} with {num_rows} rows.")

# 2. Initialize DuckDB
con = duckdb.connect(database=':memory:', read_only=False)

# 3. Query Parquet directly
print("\n--- Querying Parquet file directly with DuckDB ---")
query = f"""
SELECT
    category,
    COUNT(*) AS count,
    AVG(value) AS avg_value,
    MAX(timestamp) AS max_timestamp
FROM '{large_parquet_file}'
WHERE timestamp >= '2023-01-01' AND is_active = TRUE
GROUP BY category
ORDER BY count DESC;
"""

start_time = time.time()
result_df = con.execute(query).fetchdf()
end_time = time.time()

print(f"Query executed in {end_time - start_time:.4f} seconds.")
print(result_df)

# 4. Demonstrate Arrow integration: Create an Arrow Table and query it
print("\n--- Querying an in-memory PyArrow Table with DuckDB ---")
# Create a smaller Arrow table for in-memory demo
arrow_data = pa.table({
    'city': pa.array(['NYC', 'LA', 'Chicago', 'NYC', 'LA']),
    'population': pa.array([8.4, 3.9, 2.7, 8.4, 3.9]),
    'area_sq_mi': pa.array([302, 469, 234, 302, 469])
})

# Register the Arrow table as a view in DuckDB
con.register('my_arrow_table', arrow_data)

arrow_query = """
SELECT
    city,
    SUM(population) AS total_population
FROM my_arrow_table
WHERE area_sq_mi > 300
GROUP BY city
ORDER BY total_population DESC;
"""

start_time = time.time()
arrow_result_df = con.execute(arrow_query).fetchdf()
end_time = time.time()

print(f"Arrow Table query executed in {end_time - start_time:.4f} seconds.")
print(arrow_result_df)

# 5. Clean up
con.close()
os.remove(large_parquet_file)

アーキテクチャの比較

機能Apache ParquetApache ArrowDuckDB
目的ディスクストレージ形式インメモリデータ形式インプロセスOLAPデータベースエンジン
データ所在地永続ストレージ(HDFS、S3、ローカルディスク)RAMRAM(ディスクにスピル可能)
主な用途長期アーカイブ、データレイク、データウェアハウジングプロセス間通信、ゼロコピーデータ交換分析クエリ、ETL、データ探索
メモリレイアウトカラム型、圧縮、ページ指向カラム型、連続、非圧縮(ほとんどの場合)カラム型、ベクトル化、内部表現
シリアライズ高(ディスクI/O、圧縮/解凍)ゼロコピー(すでにArrow形式の場合)最小限(カラム型データで直接動作)
クエリエンジンなし(ストレージのみ)なし(フォーマットのみ)完全なSQLエンジン、ベクトル化実行
述語プッシュダウンあり(カラム統計経由)N/A(インメモリ、直接アクセス)あり(Parquet統計、内部最適化を活用)
SIMD最適化間接的(解凍後処理)直接的(メモリレイアウトがSIMDを可能にする)直接的(ベクトル化エンジンの核)
典型的なレイテンシ数秒から数分(大規模スキャン)数ミリ秒(データ転送)数ミリ秒から数秒(複雑なクエリ)
Advertisement

本番環境での注意点とトラブルシューティング

  1. Parquet行グループサイズの誤設定:
    • 症状: Parquetファイルでのクエリ、特にWHERE句を含むクエリが極端に遅い。
    • 問題: 小さすぎる行グループが多すぎると、メタデータの読み取りが過剰になり、オーバーヘッドが高くなります。大きすぎる行グループが少なすぎると、プルーニングの有効性が低下し、エンジンが必要以上に多くのデータを読み込むことになります。
    • 解決策: 非圧縮サイズで128MBから512MBの行グループを目指します。これは通常、スキーマの幅に応じて、行グループあたり100,000から1,000,000行に相当します。pq.write_table(..., row_group_size=...)または他の言語の同等のものを使用してください。
  2. Parquet統計の欠落または陳腐化:
    • 症状: クエリエンジン(DuckDB、Sparkなど)が、非常に選択的な述語がある場合でもフルスキャンを実行する。
    • 問題: Parquetファイルが統計なしで書き込まれた(write_statistics=False)か、データの破損/手動変更により統計が古くなっている。
    • 解決策: Parquetを書き込む際には、常にwrite_statistics=Trueを確保してください。ファイルがすでに書き込まれている場合は、書き換えを検討するか、統計を生成/更新できるツール(例:一部のデータウェアハウスのANALYZE TABLE、または再取り込み)を使用してください。
  3. 長時間実行プロセスでのArrowメモリリーク:
    • 症状: pyarrowを使用するPythonプロセスが時間の経過とともにRAM消費量を増やし、最終的にクラッシュする。
    • 問題: pyarrowオブジェクト、特にTableとRecordBatchは、オフヒープメモリを管理します。参照が意図せず保持されたり、IPCリーダー/ライターでclose()メソッドが呼び出されなかったりすると、メモリがすぐに解放されない可能性があります。
    • 解決策: 不要になった大きなArrowオブジェクト(del my_arrow_table)を明示的に削除します。IPCストリームの場合、reader.close()とwriter.close()が呼び出されていることを確認します。Pythonで積極的なガベージコレクションを行うにはgc.collect()を使用しますが、これはより深い参照管理の問題に対する一時的な対処法であることがよくあります。memrayやtracemallocなどのツールでメモリ使用量をプロファイルします。
  4. DuckDBのメモリ不足(OOM)エラー:
    • 症状: DuckDBが、複雑なクエリ、特に大規模なデータセットでの結合や集計中にクラッシュしたり、OOMエラーを報告したりする。
    • 問題: DuckDBは効率的ですが、中間結果のためにRAMが必要です。ワーキングセットが利用可能なメモリを超えると、ディスクにスピルしますが、スピル自体が大きすぎる場合やシステムがディスク容量を使い果たした場合、OOMが発生する可能性があります。
    • 解決策:
      • プロセスに利用可能なRAMを増やします。
      • クエリを最適化します: 早期にフィルタリングし、必要なカラムのみを射影します。
      • DuckDBのメモリ制限を調整します: SET memory_limit='8GB';。
      • スピル用に十分な一時ディスクスペースを確保します。
      • 非常に大規模な結合の場合は、可能であればデータを事前集計するか、異なる結合戦略を検討します。
  5. 小さなDuckDBバッチでのパフォーマンス低下:
    • 症状: DuckDBクエリが予想よりも遅い、特に小さなデータセットや最適化されていないソースからデータを処理する場合。
    • 問題: DuckDBはベクトル化されていますが、入力データソースが非常に小さなバッチ(例:カスタムイテレータから一度に1行)を提供する場合、ベクトル化処理のオーバーヘッドがメリットを上回る可能性があります。
    • 解決策: カスタムデータソースやイテレータを使用する場合は、DuckDBに適切なサイズのバッチ(例:1024〜4096行)でデータが供給されるようにします。ParquetまたはArrowを読み取る場合、DuckDBは内部でバッチ処理を処理します。

よくある質問

  1. ParquetとArrowはいつ使い分けるべきですか? Parquetは、コストとI/O効率のために圧縮と述語プッシュダウンが重要となる、データレイクやウェアハウスでのディスク上の長期的な永続ストレージに使用します。Arrowは、ゼロコピーセマンティクスと高性能なベクトル化操作が最重要となる、異なるシステム間または同じシステム内のコンポーネント間のインメモリデータ転送と処理に使用します。これらは補完的であり、排他的ではありません。

  2. DuckDBはS3や他のクラウドストレージから直接データを読み取ることができますか? はい、DuckDBはS3、GCS、Azure Blob StorageからParquetファイルとCSVファイルを直接読み取るための組み込みサポートを備えています。通常、SETコマンドまたは環境変数を使用して認証情報(例:AWSアクセスキー)を設定する必要があります。例:INSTALL httpfs; LOAD httpfs; SET s3_region='us-east-1'; SELECT * FROM 's3://my-bucket/data.parquet';。

  3. DuckDBはParquetからのデータ型とスキーマ進化をどのように処理しますか? DuckDBは、Parquetファイルのメタデータから直接スキーマを推論します。ネストされた型を含む幅広いデータ型をサポートしています。スキーマ進化の場合、新しいカラムがParquetファイルに追加されても、明示的にクエリされない限り、DuckDBは通常それらを無視します。カラムの型が互換性のない方法で変更された場合、エラーが発生します。追加的なスキーマ変更に対しては一般的に堅牢です。

  4. pyarrow.Table.to_pandas()を使用することのパフォーマンスへの影響は何ですか? 大規模なpyarrow.Tableをpandas.DataFrameに変換するには、Arrowの連続したメモリバッファからpandasの内部の、しばしば断片化されたメモリレイアウトにデータをコピーする必要があります。これはコピー操作であり、ゼロコピーではないため、非常に大規模なデータセットではパフォーマンスのボトルネックとなり、メモリを大量に消費する可能性があります。分析タスクの場合、pyarrow.Tableに対してpolarsやduckdbなどのライブラリを使用して直接操作する方が効率的です。これらのライブラリはArrowテーブルを直接消費できます。

  5. DuckDBは本番OLTPワークロードに適していますか? いいえ。DuckDBはOLAP(Online Analytical Processing)データベースであり、大規模なデータセットに対する複雑な分析クエリに最適化されています。頻繁な小さな書き込み、更新、削除を必要とする高並行性、低レイテンシのトランザクションワークロード(OLTP)向けには設計されていません。OLTPには、PostgreSQLやMySQLのような従来のリレーショナルデータベースがより適しています。DuckDBは、多くの場合SparkやPandasをローカルデータ処理に置き換える組み込み分析エンジンとして優れています。

Share this article:

Stay Updated

Get the latest posts delivered straight to your inbox.

Free Developer Utilities

Free In-Browser Developer Tools

Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.

Explore Tools
Advertisement