SingleStore 7 — Pipelines & nạp dữ liệu
SingleStore 7 — Pipelines & nạp dữ liệu
Ở các bài trước ta đã hiểu SingleStore lưu dữ liệu thế nào (rowstore/columnstore, Universal Storage) và chia dữ liệu ra sao (SHARD KEY → partition trải trên leaf nodes). Câu hỏi thực chiến tiếp theo: làm sao đưa hàng chục nghìn bản ghi mỗi giây vào cluster mà không nghẽn ở aggregator, không mất và không nhân đôi dữ liệu?
Câu trả lời đặc trưng của SingleStore là Pipelines — một cơ chế nạp dữ liệu streaming gốc, chạy ngay trong engine, kéo dữ liệu song song thẳng vào từng partition trên leaf, với đảm bảo exactly-once. Đây là điểm khác biệt lớn so với việc tự viết một consumer bên ngoài rồi bắn INSERT qua aggregator.
Bối cảnh xuyên suốt: một ngân hàng cần ingest luồng giao dịch (transactions) từ core banking, đẩy qua Kafka, vào một bảng phân tích columnstore để dashboard giám sát gian lận và báo cáo đọc gần real-time.
Lưu ý: các block SQL trong bài là cú pháp SingleStore (tương thích MySQL), dùng để minh hoạ. Sandbox của Knowledge Base là PostgreSQL read-only nên các block này không được đánh dấu "Chạy được".
Mô hình tinh thần: pipeline là consumer chạy trong cluster
Cách nạp "ngây thơ" là dựng một ứng dụng ngoài: đọc Kafka → gom lô → gửi INSERT/LOAD DATA tới aggregator. Mọi byte đi qua một điểm (aggregator), rồi aggregator mới định tuyến (route) từng dòng xuống đúng partition. Aggregator trở thành nút cổ chai, và bạn phải tự lo chuyện offset, retry, tránh nạp trùng.
Pipeline lật ngược mô hình đó. Bạn khai báo nguồn bằng DDL; cluster tự phân công mỗi partition đọc một phần của nguồn (ví dụ một tập Kafka partition), rồi nạp trực tiếp vào chính partition đó trên leaf. Không dồn qua một điểm. Master Aggregator chỉ điều phối và theo dõi tiến độ (offset/metadata), không phải kênh dẫn dữ liệu.
Ba hệ quả quan trọng của mô hình này:
- Song song thực sự & mở rộng tuyến tính: thêm leaf/partition → thêm luồng đọc song song. Thông lượng ingest scale cùng cluster.
- Không nghẽn aggregator: dữ liệu không đi qua một điểm; aggregator chỉ giữ metadata tiến độ.
- Đảm bảo ngữ nghĩa: engine tự lưu offset đã nạp trong metadata, kèm cơ chế transaction nội bộ để đạt exactly-once.
Nguồn được hỗ trợ
Pipeline nạp streaming gốc từ các nguồn sau:
| Nhóm | Nguồn |
|---|---|
| Message queue | Kafka (và các dịch vụ tương thích giao thức Kafka) |
| Object storage | Amazon S3, Azure Blob, Google Cloud Storage (GCS) |
| Distributed FS | HDFS |
| Local/mount | Filesystem |
Định dạng dữ liệu hỗ trợ gồm CSV (delimited), JSON, Avro, Parquet. Với object store (S3/Blob/GCS), pipeline theo dõi các object mới khớp tiền tố và nạp dần — biến một thư mục S3 đang được ghi liên tục thành một nguồn "gần streaming".
CREATE PIPELINE ... INTO TABLE
Dạng cơ bản: nạp thẳng từ nguồn vào một bảng. Ví dụ nạp luồng giao dịch từ Kafka vào bảng columnstore:
-- Bảng đích (columnstore, shard theo tài khoản)
CREATE TABLE transactions (
txn_id BIGINT NOT NULL,
account_no VARCHAR(20) NOT NULL,
amount DECIMAL(18,2) NOT NULL,
currency CHAR(3) NOT NULL,
kind VARCHAR(10) NOT NULL, -- 'debit' | 'credit'
txn_time DATETIME(6) NOT NULL,
SHARD KEY (account_no),
SORT KEY (txn_time),
KEY (txn_id) USING CLUSTERED COLUMNSTORE
);
-- Pipeline đọc topic Kafka, map cột theo thứ tự trường JSON
CREATE PIPELINE tx_ingest
AS LOAD DATA KAFKA 'kafka-broker-1:9092,kafka-broker-2:9092/core.transactions'
INTO TABLE transactions
FORMAT JSON (
txn_id <- txn_id,
account_no <- account_no,
amount <- amount,
currency <- currency,
kind <- kind,
txn_time <- txn_time
);
-- Kiểm thử offline: nạp thử một lô rồi xem, KHÔNG lưu offset
TEST PIPELINE tx_ingest LIMIT 5;
-- Bật pipeline chạy nền
START PIPELINE tx_ingest;
Vài điểm cần nhớ:
CREATE PIPELINEchỉ định nghĩa; pipeline chỉ chạy sauSTART PIPELINE. CóSTOP PIPELINEđể tạm dừng (offset được giữ) vàDROP PIPELINE.TEST PIPELINEcho phép xem trước dữ liệu sẽ nạp mà không commit offset — rất hợp để debug mapping.- Mỗi lô (batch) mà pipeline nạp là một giao dịch: hoặc toàn bộ lô vào bảng và offset tiến lên, hoặc không gì cả — nền tảng của exactly-once (xem phần dưới).
CREATE PIPELINE ... INTO PROCEDURE — transform khi nạp
Khi cần biến đổi dữ liệu trước khi ghi — làm sạch, phân loại, tách/gộp bảng, tra cứu bảng tham chiếu, hay ghi vào nhiều bảng — ta cho pipeline nạp vào một stored procedure thay vì thẳng vào bảng. Engine đưa cả lô dữ liệu vào procedure dưới dạng một query-typed variable (một tập hàng), và bạn xử lý bằng SQL tập hợp (set-based), không lặp từng dòng.
-- Procedure nhận nguyên một lô (batch) dữ liệu thô của pipeline
DELIMITER //
CREATE OR REPLACE PROCEDURE load_transactions(batch QUERY(
txn_id BIGINT,
account_no VARCHAR(20),
amount DECIMAL(18,2),
currency CHAR(3),
kind VARCHAR(10),
txn_time DATETIME(6)
))
AS
BEGIN
-- 1) Ghi các giao dịch hợp lệ, chuẩn hoá 'kind' về chữ thường
INSERT INTO transactions (txn_id, account_no, amount, currency, kind, txn_time)
SELECT txn_id, account_no, amount, currency, LOWER(kind), txn_time
FROM batch
WHERE account_no IS NOT NULL AND amount >= 0;
-- 2) Tách nhánh: giao dịch giá trị lớn → bảng cảnh báo (set-based)
INSERT INTO high_value_alerts (txn_id, account_no, amount, txn_time)
SELECT txn_id, account_no, amount, txn_time
FROM batch
WHERE amount > 500000000; -- > 500 triệu VND (minh hoạ)
-- 3) Bản ghi lỗi → bảng rejects để đối soát
INSERT INTO txn_rejects (txn_id, reason, raw_time)
SELECT txn_id, 'invalid_account_or_amount', txn_time
FROM batch
WHERE account_no IS NULL OR amount < 0;
END //
DELIMITER ;
CREATE PIPELINE tx_ingest_proc
AS LOAD DATA KAFKA 'kafka-broker-1:9092/core.transactions'
INTO PROCEDURE load_transactions
FORMAT JSON (
txn_id <- txn_id, account_no <- account_no, amount <- amount,
currency <- currency, kind <- kind, txn_time <- txn_time
);
START PIPELINE tx_ingest_proc;
Điểm mấu chốt: toàn bộ procedure chạy trong cùng giao dịch của lô. Nếu procedure lỗi, cả lô bị rollback và offset không tiến lên — lần sau nạp lại đúng lô đó. Nhờ vậy exactly-once vẫn giữ ngay cả khi có transform phức tạp và ghi nhiều bảng.
Exactly-once: cơ chế nào bảo chứng?
"Exactly-once" ở đây nghĩa là mỗi bản ghi nguồn được phản ánh đúng một lần vào bảng đích, kể cả khi pipeline bị dừng, leaf failover hay cluster khởi động lại. Cơ chế:
- Engine chia nguồn thành các lô (batch), mỗi lô ứng với một khoảng offset xác định.
- Việc ghi dữ liệu của lô và việc cập nhật offset đã nạp nằm trong cùng một giao dịch (atomic). Không có chuyện ghi xong dữ liệu nhưng "quên" ghi offset, hay ngược lại.
- Nếu lô thất bại giữa chừng → rollback, offset giữ nguyên → lô được nạp lại nguyên vẹn, không tạo bản ghi thừa.
Đây là lý do exactly-once của pipeline dựa trực tiếp vào mô hình giao dịch của SingleStore — chi tiết ACID, durability qua transaction log/snapshot và failover xem ở bài 8 — Giao dịch, HA & DR. (Lưu ý: với nguồn Kafka, đảm bảo này gắn với offset của Kafka; nếu procedure có side-effect ra ngoài cluster thì không nằm trong phạm vi bảo chứng.)
Pipeline vs LOAD DATA vs INSERT
Ba cách đưa dữ liệu vào, cho ba mục đích khác nhau:
| Tiêu chí | INSERT | LOAD DATA | CREATE PIPELINE |
|---|---|---|---|
| Kiểu | OLTP, từng câu | Batch một lần | Streaming liên tục |
| Đường đi | Qua aggregator → route | Qua aggregator | Song song thẳng vào partition/leaf |
| Nguồn | Giá trị trong câu lệnh | File local/nguồn | Kafka, S3, HDFS, Blob, GCS, FS |
| Theo dõi offset | Không | Không | Có (tự động) |
| Exactly-once | Không (tự lo) | Không | Có |
| Hợp cho | Ghi lẻ, cập nhật | Nạp một lần, migrate | Ingest liên tục quy mô lớn |
Quy tắc chọn: ingest liên tục từ nguồn ngoài → Pipeline. Nạp một lần một file/tập file → LOAD DATA (bản thân pipeline cũng dựa trên LOAD DATA mở rộng). Ghi giao dịch lẻ từ ứng dụng → INSERT. Tránh dùng vòng lặp INSERT từng dòng để nạp khối lượng lớn: mọi dòng dồn qua aggregator, chậm và không có exactly-once.
Giám sát pipeline
Trạng thái và tiến độ pipeline nằm trong các bảng thuộc information_schema:
-- Danh sách pipeline và trạng thái (RUNNING/STOPPED/ERROR)
SELECT DATABASE_NAME, PIPELINE_NAME, STATE
FROM information_schema.PIPELINES;
-- Lịch sử từng lô: thời gian, số dòng, số byte, lỗi (nếu có)
SELECT PIPELINE_NAME, BATCH_ID, BATCH_STATE,
BATCH_ROWS_WRITTEN, BATCH_TIME
FROM information_schema.PIPELINES_BATCHES_SUMMARY
WHERE PIPELINE_NAME = 'tx_ingest'
ORDER BY BATCH_ID DESC
LIMIT 20;
-- Lỗi nạp chi tiết để đối soát
SELECT PIPELINE_NAME, ERROR_MESSAGE, ERROR_TYPE, LOAD_DATA_LINE
FROM information_schema.PIPELINES_ERRORS
WHERE PIPELINE_NAME = 'tx_ingest'
ORDER BY ERROR_UNIX_TIMESTAMP DESC;
Các chỉ số cần theo dõi khi vận hành: độ trễ (lag) giữa offset nguồn và offset đã nạp, số lô lỗi, và thông lượng dòng/giây. Kết hợp với management views MV_* (bài 10) để cảnh báo khi ingest tụt lại sau nguồn.
Use case thực tế
Ngân hàng đẩy toàn bộ giao dịch core banking sang topic Kafka core.transactions (24 Kafka partition). Cluster SingleStore có 4 leaf, mỗi leaf 6 partition = 24 database partition — khớp 1–1 với Kafka partition.
- Trước đây: một dịch vụ Java consume Kafka rồi bắn
INSERTqua aggregator, đỉnh ~8.000 dòng/giây thì aggregator CPU chạm trần, dashboard trễ 3–5 phút (số minh hoạ). - Sau khi chuyển sang
CREATE PIPELINE ... INTO PROCEDURE: mỗi partition tự đọc phần Kafka của mình, transform (chuẩn hoá, tách bảng cảnh báo, lọc rác) ngay trong cluster. Thông lượng đạt ~35.000 dòng/giây, độ trễ dashboard xuống dưới 10 giây (minh hoạ). Khi một leaf restart, replica lên thay, pipeline tiếp tục từ đúng offset đã commit — không mất, không trùng giao dịch.
Kết quả: đội vận hành bỏ được cả một service ingest tự viết, và bộ phận giám sát gian lận có dữ liệu gần real-time để chặn giao dịch bất thường.
Ghi nhớ
- Pipeline = cơ chế nạp streaming gốc của SingleStore từ Kafka, S3, HDFS, Azure Blob, GCS, filesystem; định nghĩa bằng DDL, chạy trong engine.
- Nạp song song thẳng vào từng partition trên leaf, không dồn qua aggregator → thông lượng scale cùng cluster.
- Đảm bảo exactly-once: ghi dữ liệu lô + cập nhật offset nằm trong cùng một giao dịch; lô lỗi thì rollback và nạp lại nguyên vẹn.
INTO TABLEnạp thẳng;INTO PROCEDUREcho phép transform set-based cả lô, ghi nhiều bảng, vẫn giữ exactly-once.- Phân biệt: Pipeline cho ingest liên tục,
LOAD DATAcho nạp batch một lần,INSERTcho ghi lẻ; đừng dùngINSERTvòng lặp để nạp khối lượng lớn. - Vòng đời:
CREATE→TEST→START/STOP/DROP;TEST PIPELINExem trước không commit offset. - Giám sát qua
information_schema.PIPELINES,PIPELINES_BATCHES_SUMMARY,PIPELINES_ERRORS; theo dõi lag, lô lỗi, thông lượng.
Nguồn tham khảo
- SingleStore Documentation — "Load Data with Pipelines" / "Pipelines Overview": https://docs.singlestore.com/
- SingleStore Documentation — "CREATE PIPELINE" (cú pháp
INTO TABLE,INTO PROCEDURE,LOAD DATA KAFKA/S3/HDFS). - SingleStore Documentation — "Pipelines and Exactly-Once Semantics".
- SingleStore Documentation — "LOAD DATA".
- SingleStore Documentation — "Monitoring Pipelines" và
information_schema(PIPELINES, PIPELINES_BATCHES_SUMMARY, PIPELINES_ERRORS). - SingleStore Documentation — "Cluster Architecture" (leaf, partitions).
- MySQL Documentation — cú pháp
LOAD DATA(phần tương thích): https://dev.mysql.com/doc/
Bài viết liên quan
Điểm khác biệt lớn nhất của SingleStore: nó BIÊN DỊCH truy vấn ra mã máy (code generation) rồi chạy song song MPP trên leaf, thay vì diễn giải từng dòng. Bài dựng luồng SQL → tối ưu → sinh mã → plan biên dịch, giải thích plan cache tái dùng (biên dịch 1 lần), query pushdown xuống leaf và aggregator gộp; cách đọc EXPLAIN/PROFILE (thời gian, rows, bộ nhớ, network/reshuffle) và SHOW PLANCACHE để nhận diện reshuffle/broadcast, tối ưu truy vấn.
Kiến trúc shared-nothing của SingleStore: Master Aggregator giữ metadata và điều phối, Child Aggregator scale kết nối, Leaf node chứa dữ liệu chia thành partition. Bài mổ xẻ luồng một query (aggregator nhận → pushdown xuống leaf → gộp kết quả) và cơ chế High Availability master/replica, failover, redundancy level.
SingleStore (tiền thân MemSQL) là database quan hệ phân tán HTAP, tương thích giao thức MySQL, gộp OLTP và OLAP trong một hệ thống. Bài mở màn dựng mô hình tinh thần về HTAP, giải thích vấn đề nó giải quyết (tránh ETL sang warehouse riêng), chỉ rõ khi nào NÊN và KHÔNG NÊN dùng, định vị so với PostgreSQL, ClickHouse, BigQuery, TiDB/CockroachDB, và vẽ bản đồ toàn series 10 bài.
Khoá chính/ngoại/tổng hợp, ràng buộc (NOT NULL, UNIQUE, CHECK, FK) và cách mô hình hoá quan hệ 1:1, 1:n, n:n cho hệ khách hàng — tài khoản — giao dịch. Đi qua chuẩn hoá 1NF/2NF/3NF bằng ví dụ trước/sau cụ thể, rồi bàn khi nào nên cố tình phi chuẩn hoá để đọc nhanh — giúp thiết kế lược đồ đúng ngay từ đầu.
Cảm nhận của bạn
Bình luận
Chưa có bình luận. Hãy là người đầu tiên chia sẻ!