•17 min read

Dữ liệu cột vào năm 2026: Bố cục bộ nhớ Apache Arrow, lưu trữ Parquet & Vector hóa DuckDB

Dữ liệu cột vào năm 2026: Bố cục bộ nhớ Apache Arrow, lưu trữ Parquet & Vector hóa DuckDB

Các định dạng dữ liệu dạng cột là nền tảng cho các hệ thống phân tích hiện đại, cho phép hiệu suất truy vấn dưới một giây trên các tập dữ liệu quy mô petabyte. Hướng dẫn này phân tích sự tương tác giữa Apache Parquet để lưu trữ bền vững, Apache Arrow để xử lý trong bộ nhớ và công cụ thực thi vector hóa của DuckDB, minh họa cách các công nghệ này cùng nhau tạo thành một ngăn xếp dữ liệu phân tích hiệu suất cao.

Audio Briefing
0:00 / 0:00

Apache Parquet: Lưu trữ dạng cột tối ưu hóa cho đĩa

Apache Parquet là tiêu chuẩn thực tế để lưu trữ dữ liệu dạng cột trong các hệ sinh thái dữ liệu lớn. Thiết kế của nó ưu tiên I/O đĩa hiệu quả và predicate pushdown, những yếu tố quan trọng cho các khối lượng công việc phân tích.

Cấu trúc tệp Parquet

Một tệp Parquet bao gồm các nhóm hàng (row groups), là các phân vùng ngang của dữ liệu. Trong mỗi nhóm hàng, dữ liệu được lưu trữ theo cột. Cấu trúc này cho phép chỉ đọc các cột cần thiết cho một truy vấn, giảm đáng kể 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
│   └── ...
└── ...

Mã hóa và nén

Parquet sử dụng nhiều lược đồ mã hóa và thuật toán nén khác nhau để giảm thiểu dung lượng lưu trữ và cải thiện hiệu suất đọc.

  • Mã hóa từ điển (Dictionary Encoding): Đối với các cột có tính đa dạng thấp (low cardinality), các giá trị duy nhất được lưu trữ trong một từ điển, và dữ liệu cột lưu trữ các chỉ mục đến từ điển này. Điều này rất hiệu quả cho dữ liệu chuỗi hoặc phân loại.
  • Mã hóa độ dài chạy (Run Length Encoding - RLE): Hiệu quả cho các chuỗi giá trị lặp lại.
  • Mã hóa đóng gói bit (Bit-Packed Encoding): Đối với số nguyên, đóng gói các giá trị vào số bit tối thiểu cần thiết.
  • Nén (Compression): Snappy, Gzip, Zstd, Brotli. Zstd cung cấp sự cân bằng tuyệt vời giữa tỷ lệ nén và tốc độ.

Predicate Pushdown và Thống kê cột

Định dạng tệp của Parquet bao gồm các thống kê (giá trị min/max, số lượng null) cho mỗi khối cột (column chunk) trong một nhóm hàng. Các công cụ truy vấn tận dụng các thống kê này để loại bỏ toàn bộ nhóm hàng hoặc trang (pages) không thể thỏa mãn các điều kiện của truy vấn, tránh I/O không cần thiết.

Hãy xem xét một truy vấn SELECT sum(value) FROM table WHERE timestamp > '2023-01-01'. Nếu thống kê cột timestamp của một nhóm hàng cho biết dấu thời gian tối đa của nó là 2022-12-31, toàn bộ nhóm hàng đó sẽ bị bỏ qua.

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: Định dạng cột trong bộ nhớ

Apache Arrow định nghĩa một bố cục bộ nhớ liên tục, không phụ thuộc ngôn ngữ cho dữ liệu dạng cột. Mục tiêu chính của nó là loại bỏ chi phí tuần tự hóa/giải tuần tự hóa khi dữ liệu di chuyển giữa các hệ thống hoặc quy trình khác nhau.

Bố cục bộ nhớ Arrow

Sức mạnh của Arrow nằm ở các bộ đệm bộ nhớ được chuẩn hóa, phẳng và liên tục của nó. Đối với một cột nhất định, tất cả các giá trị cùng loại được lưu trữ cùng nhau trong một bộ đệm duy nhất. Bố cục này rất thân thiện với bộ nhớ cache và cho phép các hoạt động SIMD (Single Instruction, Multiple Data) hiệu quả.

Đối với một cột số nguyên đơn giản: [value1, value2, value3, ..., valueN]

Đối với một cột số nguyên có thể null: [validity_bitmap_buffer], [value1, value2, value3, ..., valueN] Bitmap hợp lệ cho biết giá trị nào là null.

Đối với các kiểu có độ dài thay đổi (ví dụ: chuỗi): [offset_buffer], [data_buffer], [validity_bitmap_buffer] Bộ đệm offset lưu trữ vị trí bắt đầu và độ dài của mỗi chuỗi trong bộ đệm dữ liệu.

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]}")

Đọc không sao chép (Zero-Copy Reads)

Ưu điểm đáng kể nhất của Arrow là khả năng không sao chép của nó. Khi dữ liệu đã ở định dạng Arrow, nó có thể được truyền giữa các quy trình, ngôn ngữ hoặc thậm chí các máy khác nhau (thông qua Arrow Flight) mà không cần tuần tự hóa hoặc giải tuần tự hóa. Điều này làm giảm đáng kể chu kỳ CPU và độ trễ.

Ví dụ, một tệp Parquet có thể được đọc trực tiếp vào một Table hoặc RecordBatch của Arrow. Trình đọc Parquet ánh xạ dữ liệu cột trên đĩa sang bố cục trong bộ nhớ của Arrow, thường với sự biến đổi dữ liệu tối thiểu.

DuckDB: Công cụ truy vấn vector hóa

DuckDB là một cơ sở dữ liệu SQL OLAP trong quy trình được thiết kế cho các khối lượng công việc phân tích. Sức mạnh cốt lõi của nó nằm ở công cụ thực thi truy vấn vector hóa, hoạt động trực tiếp trên dữ liệu dạng cột, thường ở định dạng Arrow, tận dụng các lệnh SIMD để đạt hiệu suất cực cao.

Thực thi vector hóa

Thay vì xử lý từng hàng một (xử lý từng tuple một), DuckDB xử lý dữ liệu theo lô (vectors) thường từ 1024-4096 hàng. Mỗi hoạt động (ví dụ: lọc, nối, tổng hợp) hoạt động trên toàn bộ các vector, cho phép:

  1. Giảm chi phí gọi hàm: Một lần gọi hàm xử lý nhiều giá trị.
  2. Hiệu quả bộ nhớ cache CPU: Dữ liệu được truy cập tuần tự, dẫn đến tỷ lệ truy cập cache cao.
  3. Tận dụng lệnh SIMD: Các CPU hiện đại có thể thực hiện cùng một thao tác trên nhiều điểm dữ liệu đồng thời bằng cách sử dụng các lệnh SIMD (ví dụ: AVX2, AVX-512). Bố cục dạng cột của DuckDB hoàn toàn phù hợp cho điều này.

Tích hợp với Parquet và Arrow

DuckDB có thể trực tiếp truy vấn các tệp Parquet mà không yêu cầu câu lệnh CREATE TABLE rõ ràng hoặc nhập dữ liệu. Nó đọc các tệp Parquet, chuyển đổi chúng thành biểu diễn cột nội bộ của nó (tương thích với Arrow), và sau đó thực thi các truy vấn bằng công cụ vector hóa của nó. Điều này loại bỏ bước ETL cho nhiều tác vụ phân tích.

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)

So sánh kiến trúc

Tính năngApache ParquetApache ArrowDuckDB
Mục đíchĐịnh dạng lưu trữ trên đĩaĐịnh dạng dữ liệu trong bộ nhớCông cụ cơ sở dữ liệu OLAP trong quy trình
Vị trí dữ liệuLưu trữ bền vững (HDFS, S3, đĩa cục bộ)RAMRAM (có thể tràn ra đĩa)
Sử dụng chínhLưu trữ dài hạn, data lakes, kho dữ liệuGiao tiếp giữa các quy trình, trao đổi dữ liệu không sao chépTruy vấn phân tích, ETL, khám phá dữ liệu
Bố cục bộ nhớDạng cột, nén, theo trangDạng cột, liên tục, không nén (chủ yếu)Dạng cột, vector hóa, biểu diễn nội bộ
Tuần tự hóaCao (I/O đĩa, nén/giải nén)Không sao chép (nếu đã ở định dạng Arrow)Tối thiểu (hoạt động trực tiếp trên dữ liệu dạng cột)
Công cụ truy vấnKhông (chỉ lưu trữ)Không (chỉ định dạng)Công cụ SQL đầy đủ, thực thi vector hóa
Predicate PushdownCó (thông qua thống kê cột)N/A (trong bộ nhớ, truy cập trực tiếp)Có (tận dụng thống kê Parquet, tối ưu hóa nội bộ)
Tối ưu hóa SIMDGián tiếp (giải nén, sau đó xử lý)Trực tiếp (bố cục bộ nhớ cho phép SIMD)Trực tiếp (cốt lõi của công cụ vector hóa)
Độ trễ điển hìnhVài giây đến vài phút (đối với các lần quét lớn)Vài mili giây (truyền dữ liệu)Vài mili giây đến vài giây (đối với các truy vấn phức tạp)
Advertisement

Những vấn đề thường gặp trong sản xuất & Khắc phục sự cố

  1. Cấu hình sai kích thước nhóm hàng Parquet:
    • Triệu chứng: Truy vấn cực kỳ chậm trên các tệp Parquet, đặc biệt với các mệnh đề WHERE.
    • Vấn đề: Quá nhiều nhóm hàng nhỏ dẫn đến đọc siêu dữ liệu quá mức và chi phí cao. Quá ít nhóm hàng lớn làm giảm hiệu quả cắt tỉa, buộc công cụ phải đọc nhiều dữ liệu hơn mức cần thiết.
    • Khắc phục: Đặt mục tiêu nhóm hàng có kích thước từ 128MB đến 512MB (chưa nén). Điều này thường tương đương với 100.000 đến 1.000.000 hàng mỗi nhóm hàng, tùy thuộc vào độ rộng của lược đồ. Sử dụng pq.write_table(..., row_group_size=...) hoặc tương đương trong các ngôn ngữ khác.
  2. Thiếu hoặc thống kê Parquet lỗi thời:
    • Triệu chứng: Các công cụ truy vấn (như DuckDB, Spark) thực hiện quét toàn bộ ngay cả với các điều kiện chọn lọc cao.
    • Vấn đề: Các tệp Parquet được ghi mà không có thống kê (write_statistics=False) hoặc thống kê đã lỗi thời do hỏng dữ liệu/sửa đổi thủ công.
    • Khắc phục: Luôn đảm bảo write_statistics=True khi ghi Parquet. Nếu các tệp đã được ghi, hãy xem xét ghi lại chúng hoặc sử dụng các công cụ có thể tạo/cập nhật thống kê (ví dụ: ANALYZE TABLE trong một số kho dữ liệu, hoặc nhập lại).
  3. Rò rỉ bộ nhớ Arrow trong các quy trình chạy dài:
    • Triệu chứng: Các quy trình Python sử dụng pyarrow tiêu thụ lượng RAM ngày càng tăng theo thời gian, cuối cùng bị treo.
    • Vấn đề: Các đối tượng pyarrow, đặc biệt là Table và RecordBatch, quản lý bộ nhớ ngoài heap. Nếu các tham chiếu được giữ lại một cách vô ý, hoặc nếu các phương thức close() không được gọi trên các trình đọc/ghi IPC, bộ nhớ có thể không được giải phóng kịp thời.
    • Khắc phục: Xóa rõ ràng các đối tượng Arrow lớn (del my_arrow_table) khi không còn cần thiết. Đối với các luồng IPC, đảm bảo reader.close() và writer.close() được gọi. Sử dụng gc.collect() để thu gom rác mạnh mẽ trong Python, mặc dù đây thường là một giải pháp tạm thời cho các vấn đề quản lý tham chiếu sâu hơn. Hồ sơ sử dụng bộ nhớ bằng các công cụ như memray hoặc tracemalloc.
  4. Lỗi DuckDB hết bộ nhớ (OOM):
    • Triệu chứng: DuckDB bị treo hoặc báo cáo lỗi OOM trong các truy vấn phức tạp, đặc biệt là các phép nối hoặc tổng hợp trên các tập dữ liệu lớn.
    • Vấn đề: DuckDB, mặc dù hiệu quả, vẫn cần RAM cho các kết quả trung gian. Nếu tập hợp làm việc vượt quá bộ nhớ khả dụng, nó sẽ tràn ra đĩa, nhưng nếu bản thân việc tràn quá lớn hoặc hệ thống hết dung lượng đĩa, OOM có thể xảy ra.
    • Khắc phục:
      • Tăng RAM khả dụng cho quy trình.
      • Tối ưu hóa truy vấn: lọc sớm, chỉ chiếu các cột cần thiết.
      • Điều chỉnh giới hạn bộ nhớ DuckDB: SET memory_limit='8GB';.
      • Đảm bảo đủ không gian đĩa tạm thời cho các lần tràn.
      • Đối với các phép nối rất lớn, hãy xem xét tổng hợp trước dữ liệu hoặc sử dụng chiến lược nối khác nếu có thể.
  5. Giảm hiệu suất với các lô DuckDB nhỏ:
    • Triệu chứng: Các truy vấn DuckDB chậm hơn dự kiến, đặc biệt trên các tập dữ liệu nhỏ hơn hoặc khi xử lý dữ liệu từ các nguồn không được tối ưu hóa.
    • Vấn đề: Mặc dù DuckDB được vector hóa, nhưng nếu nguồn dữ liệu đầu vào cung cấp các lô rất nhỏ (ví dụ: 1 hàng mỗi lần từ một trình lặp tùy chỉnh), chi phí xử lý vector hóa có thể lớn hơn lợi ích.
    • Khắc phục: Đảm bảo dữ liệu được đưa vào DuckDB theo các lô có kích thước hợp lý (ví dụ: 1024-4096 hàng) nếu sử dụng các nguồn dữ liệu hoặc trình lặp tùy chỉnh. Khi đọc Parquet hoặc Arrow, DuckDB xử lý việc phân lô nội bộ.

Các câu hỏi thường gặp

  1. Khi nào tôi nên sử dụng Parquet so với Arrow? Sử dụng Parquet để lưu trữ lâu dài, bền vững trên đĩa, đặc biệt trong các data lakes hoặc kho dữ liệu, nơi nén và predicate pushdown rất quan trọng đối với chi phí và hiệu quả I/O. Sử dụng Arrow để truyền và xử lý dữ liệu trong bộ nhớ giữa các hệ thống hoặc thành phần khác nhau trong cùng một hệ thống, nơi ngữ nghĩa không sao chép và các hoạt động vector hóa hiệu suất cao là tối quan trọng. Chúng bổ sung cho nhau, không loại trừ lẫn nhau.

  2. DuckDB có thể đọc trực tiếp dữ liệu từ S3 hoặc các bộ lưu trữ đám mây khác không? Có, DuckDB có hỗ trợ tích hợp để đọc các tệp Parquet và CSV trực tiếp từ S3, GCS và Azure Blob Storage. Bạn thường cần cấu hình thông tin xác thực (ví dụ: khóa truy cập AWS) thông qua các lệnh SET hoặc biến môi trường. Ví dụ: INSTALL httpfs; LOAD httpfs; SET s3_region='us-east-1'; SELECT * FROM 's3://my-bucket/data.parquet';.

  3. DuckDB xử lý các kiểu dữ liệu và tiến hóa lược đồ từ Parquet như thế nào? DuckDB suy luận lược đồ trực tiếp từ siêu dữ liệu tệp Parquet. Nó hỗ trợ nhiều loại dữ liệu, bao gồm các kiểu lồng nhau. Đối với tiến hóa lược đồ, nếu các cột mới được thêm vào tệp Parquet, DuckDB thường sẽ bỏ qua chúng trừ khi được truy vấn rõ ràng. Nếu các kiểu cột thay đổi không tương thích, nó sẽ báo lỗi. Nó thường mạnh mẽ đối với các thay đổi lược đồ bổ sung.

  4. Ý nghĩa hiệu suất của việc sử dụng pyarrow.Table.to_pandas() là gì? Chuyển đổi một pyarrow.Table lớn thành một pandas.DataFrame liên quan đến việc sao chép dữ liệu từ các bộ đệm bộ nhớ liên tục của Arrow vào bố cục bộ nhớ nội bộ, thường bị phân mảnh của pandas. Đây là một hoạt động sao chép, không phải không sao chép, và có thể là một nút thắt cổ chai về hiệu suất và tiêu tốn bộ nhớ đối với các tập dữ liệu rất lớn. Đối với các tác vụ phân tích, thường hiệu quả hơn khi thực hiện các hoạt động trực tiếp trên pyarrow.Table bằng cách sử dụng các thư viện như polars hoặc duckdb có thể tiêu thụ trực tiếp các bảng Arrow.

  5. DuckDB có phù hợp cho các khối lượng công việc OLTP sản xuất không? Không. DuckDB là một cơ sở dữ liệu OLAP (Online Analytical Processing), được tối ưu hóa cho các truy vấn phân tích phức tạp trên các tập dữ liệu lớn. Nó không được thiết kế cho các khối lượng công việc giao dịch có độ đồng thời cao, độ trễ thấp (OLTP) yêu cầu các thao tác ghi, cập nhật và xóa nhỏ thường xuyên. Đối với OLTP, các cơ sở dữ liệu quan hệ truyền thống như PostgreSQL hoặc MySQL phù hợp hơn. DuckDB nổi trội như một công cụ phân tích nhúng, thường thay thế Spark hoặc Pandas cho việc xử lý dữ liệu cục bộ.

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