Import CSV: đo row-by-row, batch và COPY trên cùng dữ liệu
Câu hỏi: tăng batch giúp được bao nhiêu trong một phép thử cụ thể, và làm sao biết cách nhanh hơn vẫn nạp đúng dữ liệu?
Cần biết trước: Python/SQL cơ bản và lab database dùng riêng.
Chép các file của cách A vào thư mục trống. Lab dưới đây dùng PostgreSQL18.6,
Python3.14.4 và psql, chỉ thư viện chuẩn; không cần driver hay ORM. Docker/Linux,
MySQL import, dataset10M và cold cache chưa đo.
Giữ bài toán giống nhau
Một bảng mới trong namespace lab, dữ liệu đơn hàng giả. Mỗi hàng có ID, customer ID và số tiền cent. Dữ liệu dùng công thức với seed17, không lấy thông tin thật. Ba phương án giữ nguyên logged table, PK, FK, CHECK, secondary index và một transaction cho cả file. Row-by-row ở đây là một câu INSERT/hàng, không phải một commit/hàng. Vì vậy không gán chênh lệch đo được cho autocommit.
| Phương án | Công việc client | Công việc gửi database |
|---|---|---|
| row | Python đọc CSV, chuyển số và dựng SQL | N câu INSERT trong một session/transaction |
| batch | Cùng chuyển số, gom tối đa400hàng | Một INSERT nhiều VALUES/lô, cùng transaction |
| copy | Client psql đọc file CSV | Một \copy, server phân tích CSV và kiểm constraint |
\copy đọc file ở client, khác COPY FROM '/path' đọc file trên server. Cả hai
vẫn phải xử lý dữ liệu và kiểm ràng buộc. Lab không dùng FREEZE, unlogged table,
ON_ERROR ignore hoặc tắt WAL/durability. COPY,
psql và \copy.
Chốt phép đo trước khi chạy
- Các cỡ1000/5000/10000, ba lượt mỗi tổ hợp: 27 sample. Warm-up riêng không tính.
- Thứ tự row/batch/copy luân phiên qua ba lượt để bớt ưu thế của cách chạy sau.
- Timer gồm đọc/serializeCSV, khởi tạo psql, gửi SQL, thực hiện và commit; không gồm sinh CSV, reset bảng hay validation sau nạp. Cả ba đều tạo connection mới mỗi lượt.
- Reset bằng TRUNCATE giữ schema/index/constraint. Không xóa cache hệ điều hành hay shared buffers: đây là phép đo sau warm-up, không phải cold-start.
- Đối chiếu count và SHA256 của CSV export theo ID, so toàn bộ ô với input, không chỉ count hoặc tổng tiền. Validation chạy ngoài timer.
- Cỡ này nằm trong RAM; row/batch dựng request trong bộ nhớ. Không gọi đây là streaming importer có bộ nhớ giới hạn. Protocol, serialization và parser cùng đổi giữa các cách.
Không dùng phép đo này để cô lập network RTT: tất cả connection qua socket local. Gom transaction cũng đổi chi phí commit; phải đo riêng nếu muốn so autocommit. Hướng dẫn nạp dữ liệu.
Mã chạy được từ thư mục lab
Lưu bulk.py. Hàm load dùng cùng một transaction, gặp lỗi thì psql thoát và connection
đóng khiến transaction chưa commit bị rollback. Không tiếp tục chạy COMMIT sau lỗi.
from __future__ import annotations
import csv
import hashlib
import json
import os
import platform
import statistics
import subprocess
import sys
import time
from datetime import UTC, datetime
from pathlib import Path
from typing import TypedDict
class Sample(TypedDict):
n: int
method: str
repeat: int
seconds: float
csv_sha256: str
LAB_DIR = os.environ["LAB_DIR"]
METHODS = ("row", "batch", "copy")
SIZES = (1000, 5000, 10000)
SEED = 17
def sql(statement: str) -> bytes:
env = dict(os.environ)
for key in ("PGSERVICE", "PGHOSTADDR", "PGPASSWORD", "PGOPTIONS"):
env.pop(key, None)
env.update(
PGSERVICEFILE="/dev/null",
PGPASSFILE=str(Path(LAB_DIR) / "no-pgpass"),
PGOPTIONS="-c client_min_messages=warning",
)
command = [
"psql",
"-X",
"-q",
"-w",
"-h",
LAB_DIR,
"-p",
"5432",
"-U",
"lab",
"-d",
"postgres",
"-v",
"ON_ERROR_STOP=1",
"-At",
]
return subprocess.run(
command,
input=statement.encode(),
capture_output=True,
check=True,
env=env,
timeout=60,
).stdout
def generate(path: Path, n: int) -> None:
with path.open("w", newline="") as output:
writer = csv.writer(output, lineterminator="\n")
writer.writerow(("id", "customer_id", "total_cents"))
writer.writerows(
(i, i * SEED % 2000 + 1, (i * 7919 + SEED) % 50000 + 100)
for i in range(1, n + 1)
)
def load(method: str, path: Path) -> None:
if method == "copy":
body = f"\\copy wiki_lab.import_orders FROM '{path.name}' WITH (FORMAT csv, HEADER true)\n"
else:
with path.open(newline="") as source:
reader = csv.reader(source)
next(reader)
rows = list(reader)
if any(len(row) != 3 for row in rows):
raise ValueError("CSV phải có đúng ba cột")
values = ["(" + ",".join(str(int(cell)) for cell in row) + ")" for row in rows]
batch_size = 1 if method == "row" else 400
body = (
"\n".join(
"INSERT INTO wiki_lab.import_orders VALUES "
+ ",".join(values[i : i + batch_size])
+ ";"
for i in range(0, len(values), batch_size)
)
+ "\n"
)
sql("BEGIN;\n" + body + "COMMIT;\n")
def validate(path: Path, n: int) -> str:
count = int(sql("SELECT count(*) FROM wiki_lab.import_orders;"))
exported = sql(
"\\copy (SELECT id,customer_id,total_cents FROM wiki_lab.import_orders ORDER BY id) TO STDOUT WITH (FORMAT csv, HEADER true)\n"
)
expected = path.read_bytes()
assert count == n and exported == expected, "count hoặc nội dung khác CSV"
return hashlib.sha256(exported).hexdigest()
def fault_checks(valid: Path) -> int:
with valid.open(newline="") as source:
original = list(csv.reader(source))
completed = 0
for kind in ("duplicate", "foreign_key", "check", "type", "columns"):
rows = [row.copy() for row in original]
if kind == "duplicate":
rows[-1][0] = rows[1][0]
elif kind == "foreign_key":
rows[-1][1] = "9999"
elif kind == "check":
rows[-1][2] = "-1"
elif kind == "type":
rows[-1][2] = "bad_int"
else:
rows[-1].append("extra")
bad = Path("bad.csv")
with bad.open("w", newline="") as output:
csv.writer(output, lineterminator="\n").writerows(rows)
for method in METHODS:
sql("TRUNCATE wiki_lab.import_orders;")
try:
load(method, bad)
except subprocess.CalledProcessError, ValueError:
pass
else:
raise AssertionError(f"input {kind} được chấp nhận bởi {method}")
assert int(sql("SELECT count(*) FROM wiki_lab.import_orders;")) == 0
load(method, valid)
validate(valid, 1000)
completed += 1
return completed
def main() -> None:
assert Path(LAB_DIR, ".wiki-lab").is_file()
assert sql("SHOW unix_socket_directories;").decode().strip() == LAB_DIR
settings = (
sql(
"SELECT current_setting('server_version_num'),current_setting('fsync'),current_setting('full_page_writes'),current_setting('synchronous_commit'),current_setting('wal_level'),current_setting('shared_buffers');"
)
.decode()
.strip()
.split("|")
)
assert settings[:4] == ["180006", "on", "on", "on"] and settings[4] == "replica"
sql("""
DROP TABLE IF EXISTS wiki_lab.import_orders;
CREATE TABLE wiki_lab.import_orders (
id integer PRIMARY KEY,
customer_id integer NOT NULL REFERENCES wiki_lab.customers(id),
total_cents integer NOT NULL CHECK (total_cents >= 0)
);
CREATE INDEX import_customer ON wiki_lab.import_orders(customer_id,id);
""")
for n in SIZES:
generate(Path(f"data-{n}.csv"), n)
warm = Path("data-1000.csv")
for method in METHODS:
sql("TRUNCATE wiki_lab.import_orders;")
load(method, warm)
validate(warm, 1000)
samples: list[Sample] = []
for n in SIZES:
path = Path(f"data-{n}.csv")
for repeat in range(3):
order = METHODS[repeat:] + METHODS[:repeat]
for method in order:
sql("TRUNCATE wiki_lab.import_orders;")
started = time.perf_counter()
load(method, path)
elapsed = time.perf_counter() - started
digest = validate(path, n)
samples.append(
{
"n": n,
"method": method,
"repeat": repeat + 1,
"seconds": elapsed,
"csv_sha256": digest,
}
)
fault_cases = fault_checks(warm)
assert len(samples) == 27 and fault_cases == 15
summaries = []
for n in SIZES:
for method in METHODS:
durations = [
r["seconds"] for r in samples if r["n"] == n and r["method"] == method
]
assert len(durations) == 3
mean = statistics.mean(durations)
summaries.append(
{
"n": n,
"method": method,
"mean_s": mean,
"min_s": min(durations),
"max_s": max(durations),
"stdev_s": statistics.stdev(durations),
"rows_per_s": n / mean,
}
)
metadata = {
"checked_at": datetime.now(UTC).isoformat(),
"python": platform.python_version(),
"system": platform.system(),
"architecture": platform.machine(),
"logical_cpus": os.cpu_count(),
"postgresql": settings,
"seed": SEED,
"batch": 400,
"cache": "warm-up, no eviction",
"scope": "CSVread/serialize/psql/commit",
}
payload = {
"metadata": metadata,
"samples": samples,
"summaries": summaries,
"fault_cases": fault_cases,
"retries": fault_cases,
}
Path("bulk-results.json").write_text(json.dumps(payload, indent=2) + "\n")
sys.stdout.write("BENCH_RESULT " + json.dumps(payload) + "\n")
sys.stdout.write(
"samples=27 checksum=all_equal fault_cases=15 retries=15 durability=on\n"
)
sql("DROP TABLE wiki_lab.import_orders;")
if __name__ == "__main__":
main()
Dựng fixture và chạy
set -euo pipefail
. ./lab-local.sh
lab_up
lab_seed
lab_whoami
lab_versions
set -euo pipefail
. ./lab.env
python3 bulk.py
samples=27 checksum=all_equal fault_cases=15 retries=15 durability=on
bulk-results.json giữ từng sample và metadata của lần chạy. So nội dung thành công
không chứng minh toàn importer production: ở đây dữ liệu mới, không upsert hay hiệu ứng
ngoài database. Với input lỗi, cả ba phải để0hàng; sau sửa file, retry nạp đúng1000hàng.
Không bỏ constraint để nhận tốc độ đẹp hơn. Với file lớn cần staging, checkpoint và
chính sách retry phù hợp transaction boundary; đó là bài toán khác phép đo này.
Cleanup luôn phải chạy
set -euo pipefail
. ./lab-local.sh
if [ -f lab.env ]; then
. ./lab.env
old_dir=$LAB_DIR
lab_clean
[ ! -e "$old_dir" ]
fi
lab_clean
[ ! -e lab.env ]
echo 'cleanup ok'
Kết quả và giới hạn kết luận
Đo ngày2026-10-03 lúc07:21:54UTC: macOS27.0.1, arm64,10logicalCPU, Python3.14.4 và PostgreSQL18.6 qua Unix socket; shared_buffers128MB. fsync/full_page_writes/synchronous_commit đều on, wal_level=replica. Đây là một máy local, workload nhỏ sau warm-up; không đo cạnh tranh production. Thời gian là giây, throughput = số hàng / mean wall time; stddev chỉ từ ba lượt. Khi đọc kết quả, kiểm từng sample trong JSON thay vì chỉ nhìn mean.
| Hàng | Cách | Mean (s) | Min (s) | Max (s) | Stddev (s) | Hàng/s |
|---|---|---|---|---|---|---|
| 1000 | row | 0.028945 | 0.027562 | 0.029821 | 0.001212 | 34548 |
| 1000 | batch | 0.014548 | 0.013935 | 0.015769 | 0.001058 | 68738 |
| 1000 | copy | 0.014162 | 0.013496 | 0.014676 | 0.000605 | 70613 |
| 5000 | row | 0.099265 | 0.097083 | 0.100486 | 0.001894 | 50370 |
| 5000 | batch | 0.029153 | 0.028890 | 0.029435 | 0.000273 | 171511 |
| 5000 | copy | 0.024375 | 0.024322 | 0.024432 | 0.000055 | 205125 |
| 10000 | row | 0.191379 | 0.189333 | 0.194709 | 0.002908 | 52252 |
| 10000 | batch | 0.049266 | 0.047273 | 0.052186 | 0.002584 | 202981 |
| 10000 | copy | 0.037493 | 0.037410 | 0.037620 | 0.000112 | 266716 |
Ở10000hàng, COPY đạt mean0.037493s, batch0.049266s, row0.191379s. Ở1000hàng, khoảng min/max của batch và COPY chồng nhau; ba lượt chưa đủ để kết luận COPY luôn hơn batch. Khởi tạo client và serialization vẫn nằm trong timer. Tất cả27lượt khớp toàn bộ CSV export;15input lỗi để0hàng và15retry nạp đúng. Lần chạy lại tạo JSON mới; không kỳ vọng timing hoặc checksum của metadata giống nhau.
Vì sao không suy10M là đã đo?
Mô hình tỷ lệ T(10M) ≈ T(10k) × 1000 chỉ là giả định về cùng throughput. Nó bỏ qua
giới hạn RAM của request dựng sẵn, tốc độ WAL/đĩa, checkpoint, cache, index lớn,
constraint, contention và replication. Cần benchmark cỡ lớn hơn cùng workload để
kiểm từng giả định; chưa chạy10M thì không ghi một mốc thời gian thành kết quả đo.
Trong PostgreSQL, thay wal_level, bỏ index/constraint hoặc giảm bảo đảm durability
làm phép so khác điều kiện. COPY FREEZE liên quan tuple freezing, không đồng nghĩa
“không ghi WAL”. Giữ cấu hình và đo lại nếu một điều kiện đổi; đọc trade-off thay vì
copy lệnh tắt kiểm tra. Độ tin cậy WAL.
Học tiếp: chi phí index, EXPLAIN ANALYZE và replication lag. Không có ngưỡng số hàng chung buộc mọi ứng dụng bỏ ORM hay bắt buộc chọn một batch size.