DE

[DE Blog #15] Giải Mã Công Cụ Truy Vấn Phân Tán: Kiến Trúc MPP, Trino (Presto) vs. ClickHouse & Kỹ Thuật Query Federation

9 views
[DE Blog #15] Giải Mã Công Cụ Truy Vấn Phân Tán: Kiến Trúc MPP, Trino (Presto) vs. ClickHouse & Kỹ Thuật Query Federation

1. Bối cảnh thực tế (Context & Problem Statement)

Trong một doanh nghiệp có hệ thống dữ liệu phát triển, dữ liệu thường bị phân mảnh ở nhiều nơi:

  • Dữ liệu giao dịch người dùng nằm trong PostgreSQL / MySQL.
  • Dữ liệu sự kiện lịch sử hàng trăm Terabyte nằm trên S3 Data Lakehouse (Parquet/Iceberg).
  • Dữ liệu dòng thời gian thực nằm trong Apache Kafka.
  • Dữ liệu tìm kiếm / log nằm trong Elasticsearch.

Nếu mỗi khi cần phân tích, Data Engineer lại phải viết pipeline ETL để gom tất cả về một chỗ thì sẽ tốn rất nhiều thời gian, chi phí lưu trữ gấp đôi và làm trễ thông tin phục vụ kinh doanh.

Trino (trước đây là PrestoSQL)ClickHouse ra đời như những giải pháp đột phá:

  • Trino: Công cụ truy vấn phân tán MPP mạnh mẽ cho phép chạy một câu lệnh SQL duy nhất để JOIN trực tiếp bảng trên PostgreSQL với bảng trên S3 và Kafka mà không cần di chuyển dữ liệu (Data Virtualization / Query Federation).
  • ClickHouse: Đỉnh cao của tốc độ xử lý thời gian thực nhờ kỹ thuật Vectorized Execution, quét hàng tỷ dòng dữ liệu chỉ trong vài chục mili-giây.

2. Các Khái Niệm & Cơ Chế Cốt Lõi

2.1. Kiến Trúc MPP (Massively Parallel Processing) Của Trino

Trino được thiết kế theo mô hình điều phối phân tán không chia sẻ tài nguyên (Shared-Nothing MPP):

  • Trino Coordinator: Máy chủ đầu não chịu trách nhiệm:
    1. Tiếp nhận câu lệnh SQL từ client.
    2. Phân tích cú pháp (Parsing), kiểm tra ngữ nghĩa và tối ưu hóa qua Cost-Based Optimizer (CBO).
    3. Chia nhỏ kế hoạch thực thi thành các giai đoạn (Stages) và các nhiệm vụ nhỏ (Tasks).
    4. Lập lịch và phân phối Tasks xuống các Worker Nodes.
  • Trino Workers: Các máy chủ tính toán thực thi Tasks, trực tiếp kéo dữ liệu (Splits) từ các nguồn bên ngoài qua Connector, xử lý trong bộ nhớ (In-Memory Pipelining) và trao đổi dữ liệu trung gian qua mạng (Exchange Operator).
[SQL Client / BI Tool]
          |
          v
+-------------------------------------------------------------+
|                      TRINO COORDINATOR                      |
|       (Parser -> Cost-Based Optimizer -> Scheduler)         |
+------------------------------+------------------------------+
                               |
            +------------------+------------------+
            |                                     |
            v (Tasks)                             v (Tasks)
    +---------------+                     +---------------+
    | TRINO WORKER  | <=== (Exchange) ==> | TRINO WORKER  |
    | (Memory Pipe) |                     | (Memory Pipe) |
    +-------+-------+                     +-------+-------+
            |                                     |
    +-------+-----------------------------+-------+-------+
    | (SPI Connectors)                    |               |
    v                                     v               v
[PostgreSQL DB]                 [Apache Iceberg on S3] [Kafka Stream]

2.2. Cơ Chế Query Federation & Connector Architecture (SPI)

Trino sử dụng kiến trúc giao diện mở SPI (Service Provider Interface):

  • Connector: Đóng vai trò như một "Driver" dịch chuyển ngữ nghĩa SQL của Trino sang API của hệ thống nguồn (Iceberg Connector, Hive Connector, PostgreSQL Connector, Delta Lake Connector).
  • Pushdown Optimization (Đẩy tối ưu xuống nguồn):
    • Predicate Pushdown: Đẩy mệnh đề WHERE xuống nguồn để nguồn tự lọc dữ liệu trước khi trả về qua mạng.
    • Projection Pushdown: Chỉ yêu cầu nguồn đọc đúng các cột có trong SELECT.
    • Limit / Aggregate Pushdown: Đẩy các hàm COUNT(), LIMIT xuống cho RDBMS tính trước.

2.3. Kỹ Thuật Xử Lý Dữ Liệu: Volcano Iterator vs. Vectorized Execution

  • Mô hình truyền thống (Volcano Iterator Model / Row-at-a-time): Xử lý từng dòng dữ liệu một qua lời gọi hàm next(). Gây ra hàng tỷ lần gọi hàm ảo (Virtual Function Calls), làm xáo trộn CPU Instruction Cache và không tận dụng được phần cứng hiện đại.
  • Vectorized Execution (Xử lý theo khối Vector trong ClickHouse / Trino):
    • Gom dữ liệu thành từng mảng cột nhỏ (ví dụ mảng 4096 phần tử cùng kiểu dữ liệu).
    • Vừa khít với CPU L1/L2 Cache, triệt tiêu chi phí gọi hàm.
    • Tận dụng tập lệnh SIMD (Single Instruction, Multiple Data) của CPU hiện đại để tính toán song song trên nhiều giá trị cùng một chu kỳ xung nhịp.

2.4. In-Memory Pipelining vs. MapReduce/Spark Disk Spilling

  • Apache Spark: Mặc định lưu trữ dữ liệu trung gian của các bước Shuffle ra đĩa (Disk Spill) để đảm bảo khả năng chịu lỗi (Fault Tolerance) cho các tác vụ Batch chạy nhiều giờ.
  • Trino: Dữ liệu được truyền trực tiếp qua mạng giữa các Worker qua bộ nhớ RAM (Pipelined Execution). Không ghi đĩa trung gian \rightarrow Độ trễ truy vấn cực thấp (Interactive Latency). Nếu một Worker chết giữa chừng, câu query sẽ fail ngay và cần chạy lại (fail-fast model).

3. Các Bảng Markdown So Sánh Chi Tiết

📊 Bảng 1: So Sánh Toàn Diện: Trino vs. Apache Spark vs. ClickHouse

Tiêu chíTrino (Presto)Apache SparkClickHouse
Mục đích chínhInteractive SQL Analytics & Query FederationLarge-scale Batch Processing, ETL, Machine LearningReal-time Analytics, Log Processing, Time-series Data
Mô hình tính toánMPP In-Memory Streaming PipelineDistributed Map-Reduce / Directed Acyclic Graph (DAG)Vectorized Columnar Engine trên máy chủ độc lập/cluster
Lưu trữ dữ liệuKhông lưu trữ (Compute-only Query Engine)Không lưu trữ (tích hợp storage ngoài)Có tầng lưu trữ riêng tối ưu hóa cao (MergeTree Engine)
Cơ chế chịu lỗi (Fault-tolerance)Chạy lại toàn bộ query nếu có node chết (Fail-fast)Tự động thử lại từng Task/Stage bị lỗi (Rất mạnh)Hỗ trợ Replication theo cụm (ClickHouse Keeper/ZooKeeper)
Tốc độ Interactive QueryRất nhanh (vài giây)Chậm hơn (có overhead lập lịch & shuffle đĩa)Siêu tốc (< 100ms trên dữ liệu lớn)
Khả năng Query FederationXuất sắc nhất (JOIN chéo nhiều database dễ dàng)Có hỗ trợ nhưng nặng nề hơnHạn chế (chủ yếu tối ưu dữ liệu nội bộ)

📊 Bảng 2: So Sánh Volcano Model vs. Vectorized Execution Model

Tiêu chíVolcano Iterator Model (Row-at-a-time)Vectorized Execution Model (Batch/Columnar)
Đơn vị xử lýTừng dòng dữ liệu đơn lẻ (Tuple)Một khối mảng (Block/Vector) chứa hàng nghìn giá trị cùng cột
Chi phí gọi hàm CPURất cao (NN dòng =N= N lần gọi hàm next())Cực thấp (NN dòng chỉ cần N/4096N / 4096 lần gọi hàm)
Tận dụng CPU CacheKém (Dữ liệu dạng dòng phân tán trong RAM)Hoàn hảo (Dữ liệu liên tục nạp trọn trong L1/L2 Cache)
Tận dụng SIMD HardwareKhông thểTự động tận dụng tập lệnh vector của CPU (AVX-512, NEON)
Engine tiêu biểuPostgreSQL, MySQL, Hive cũClickHouse, DuckDB, Trino, Snowflake, Databricks Photon

📊 Bảng 3: So Sánh Data Virtualization (Query Federation) vs. Physical Data Ingestion (ELT)

Tiêu chíData Virtualization (Query Federation với Trino)Physical Data Ingestion (ELT vào Data Warehouse)
Độ trễ dữ liệu (Data Freshness)Tức thì (Zero Latency - Đọc trực tiếp từ nguồn)Có độ trễ (phụ thuộc vào chu kỳ chạy pipeline ELT)
Chi phí lưu trữ (Storage)Bằng 0 (Không cần nhân bản dữ liệu)Tốn kém (Phải trả tiền lưu trữ cho Data Warehouse)
Tác động lên Database nguồnCao (Query phân tích trực tiếp có thể gây chậm nguồn OLTP)Rất thấp (Chỉ đọc 1 lần qua CDC/Batch rồi thôi)
Hiệu năng JOIN bảng lớnKém hơn nếu phải truyền hàng tỷ dòng qua mạng internetTối ưu nhất (Dữ liệu đã được nạp sẵn cùng định dạng trong DWH)
Use-case tối ưuKhám phá dữ liệu (Ad-hoc Analysis), Báo cáo dữ liệu liveBáo cáo định kỳ quy mô lớn, Dashboard cho hàng nghìn người dùng

4. Best Practices & Tối Ưu Hóa Truy Vấn Phân Tán (Pro-Tips & Tuning)

💡 Cạm bẫy "Cross-Network Data Shuffling" trong Query Federation

  • Cạm bẫy: Viết câu lệnh JOIN giữa bảng Fact lớn trên S3 (100M100\text{M} dòng) với bảng Dimension trên PostgreSQL (10M10\text{M} dòng) mà không có điều kiện lọc.
    \rightarrow Trino sẽ phải kéo toàn bộ 10M10\text{M} dòng từ PostgreSQL qua kết nối mạng JDBC, gây nghẽn băng thông và làm chậm cả Database nghiệp vụ!
  • Giải pháp:
    1. Luôn sử dụng bộ lọc (WHERE) chặt chẽ trên bảng RDBMS để kích hoạt Predicate Pushdown.
    2. Bật tính năng Dynamic Partition Pruning (DPP) trong Trino để tự động đẩy các giá trị lọc từ bảng Dimension sang bảng Fact nhằm bỏ qua các partition không cần thiết trên S3.

💡 Chọn đúng chiến lược JOIN Distribution trong Trino

  • Broadcast Join (broadcast): Gửi toàn bộ bảng nhỏ sang tất cả các Worker Nodes.
    \rightarrow Tối ưu khi một trong hai bảng có kích thước nhỏ (<100MB< 100\text{MB}).
  • Distributed Hash Join (partitioned): Băm cả hai bảng theo Join Key và phân phối đều cho các Worker.
    \rightarrow Bắt buộc sử dụng khi cả hai bảng đều lớn để tránh tràn bộ nhớ Worker.
-- Ép kiểu Join trong Trino nếu CBO ước tính sai:
SET SESSION join_distribution_type = 'BROADCAST';

5. Góc Ôn Luyện Phỏng Vấn (Interview Corner)

❓ Câu hỏi 1: Tại sao Trino (Presto) lại có thể thực thi các câu truy vấn SQL tương tác (Interactive Queries) nhanh hơn đáng kể so với Apache Spark? Khi nào bạn nên chọn Trino và khi nào nên chọn Spark?

  • Gợi ý trả lời:
    • Lý do Trino nhanh hơn:
      1. Pure In-Memory Pipelining: Trino truyền dữ liệu trực tiếp giữa các công đoạn tính toán qua bộ nhớ RAM và mạng (Streaming Exchange), không có chi phí ghi dữ liệu tạm ra đĩa cứng như Spark.
      2. Luôn chạy sẵn (Always-on Long-running Cluster): Các Worker của Trino luôn ở trạng thái chờ sẵn, không mất chi phí khởi tạo JVM, yêu cầu cấp phát tài nguyên từ Cluster Manager (YARN/K8s) cho từng query như Spark.
      3. Tối ưu hóa chuyên biệt cho SQL: Toàn bộ engine được viết chuyên biệt cho SQL Analytics thay vì phải hỗ trợ cả Graph, MLlib, RDD đa dụng.
    • Nguyên tắc lựa chọn:
      • Chọn Trino khi: Cần phục vụ các truy vấn Ad-hoc tương tác nhanh (Ad-hoc exploration), làm nền tảng Query Engine cho BI Dashboard, hoặc cần Query Federation đa nguồn.
      • Chọn Spark khi: Xây dựng các pipeline ETL/ELT nặng nề chạy nhiều giờ, xử lý dữ liệu quy mô Petabyte yêu cầu khả năng chịu lỗi cao (Fine-grained Fault Tolerance) để không bị chết giữa chừng khi có node sập.

❓ Câu hỏi 2: Giải thích cơ chế hoạt động của Pushdown Optimization trong Trino/Presto. Làm thế nào nó giúp bảo vệ cơ sở dữ liệu nguồn khi thực hiện Query Federation?

  • Gợi ý trả lời:
    • Bản chất: Pushdown Optimization là kỹ thuật mà Trino Coordinator chuyển giao một phần công việc tính toán (lọc, chọn cột, tổng hợp) cho chính hệ thống lưu trữ bên dưới thực hiện thay vì kéo toàn bộ dữ liệu thô về Trino rồi mới xử lý.
    • Các dạng Pushdown chính:
      1. Predicate Pushdown: Khi người dùng viết WHERE status = 'ACTIVE' AND created_date >= '2026-08-01', Connector sẽ biên dịch điều kiện này thành câu lệnh SQL gốc gửi tới PostgreSQL/MySQL. Database nguồn chỉ quét và trả về các dòng khớp điều kiện.
      2. Projection Pushdown: Nếu bảng nguồn có 100 cột nhưng câu lệnh chỉ SELECT id, total_amount, Connector chỉ đọc đúng 2 cột, giảm 98% dung lượng truyền tải qua mạng.
    • Bảo vệ hệ thống nguồn: Giúp giảm tối đa lượng I/O đĩa và băng thông mạng truyền tải giữa Database nguồn và cụm Trino, tránh tình trạng quét toàn bộ bảng (Full Table Scan) làm sập hệ thống Production.

❓ Câu hỏi 3: Vectorized Execution là gì? Tại sao các hệ quản trị hiện đại như ClickHouse, DuckDB và Snowflake đều chuyển đổi sang mô hình này thay vì dùng Volcano Iterator Model truyền thống?

  • Gợi ý trả lời:
    • Định nghĩa: Vectorized Execution là mô hình xử lý dữ liệu theo từng khối mảng (Vectors/Batches) gồm hàng nghìn phần tử của cùng một cột tại một thời điểm, thay vì xử lý từng dòng dữ liệu tuần tự.
    • Lý do các hệ thống hiện đại áp dụng:
      1. Giảm thiểu tối đa chi phí CPU Overhead: Thay vì thực hiện hàng triệu lời gọi hàm ảo (Virtual Function Calls) cho từng dòng, hệ thống chỉ gọi một hàm duy nhất cho cả một khối mảng.
      2. Tối ưu hóa phân cấp bộ nhớ (Cache Locality): Dữ liệu của một cột có cùng kiểu dữ liệu được nạp liên tục và nằm trọn vẹn trong CPU L1/L2 Cache, loại bỏ hoàn toàn hiện tượng Cache Miss.
      3. Khai thác phần cứng song song (SIMD Instructions): Trình biên dịch có thể tự động ánh xạ các phép toán (ví dụ cộng 2 cột số: A+BA + B) thành các tập lệnh vector hóa của CPU (Intel AVX-512, ARM NEON), cho phép tính toán đồng thời 8 hoặc 16 phép toán trong cùng một chu kỳ xung nhịp.

6. Tóm Tắt Ghi Nhớ Nhanh (Key Takeaways)

  1. Trino (Presto) là công cụ hàng đầu cho Interactive Querying & Query Federation (MPP Architecture).
  2. In-Memory Streaming Pipelining giúp Trino đạt tốc độ vượt trội nhưng đánh đổi bằng việc không có cơ chế lưu vết trung gian (Fail-fast).
  3. Vectorized Execution + SIMD trong ClickHouse/DuckDB là tiêu chuẩn hiệu năng hiện đại, xử lý dữ liệu theo khối cột trong CPU Cache.
  4. Luôn tận dụng Pushdown Optimization để tránh biến Query Federation thành thảm họa nghẽn băng thông mạng.
0Claps