Chương 6: Idempotent Replay Verification & System Integration#

Thành viên phụ trách: Huỳnh Trọng Viên


1. Bối cảnh và Mục tiêu#

Chương này đảm nhận vai trò thẩm định lại toàn bộ chất lượng của đường ống dữ liệu (pipeline) mà nhóm đã xây dựng ở các bước trước. Mục tiêu trọng yếu là kiểm chứng thực nghiệm hai tính chất kiến trúc cốt lõi: tính Idempotent (Lũy đẳng) và tính Incremental (Xử lý tăng tiến).

Kịch bản thực nghiệm yêu cầu tạo ra sự biến đổi trên đúng 1 tệp mã nguồn duy nhất trong repository, kích hoạt Parser Service quét mục tiêu thay đổi cục bộ này (parse-changed), và tiến hành đối chiếu. Hệ thống được đánh giá là thành công khi và chỉ khi dữ liệu trong Neo4j và MongoDB phản ánh được sự cập nhật, không làm gia tăng kích thước tập dữ liệu (không sinh bản ghi trùng), đồng thời cơ chế Spark Checkpoint nhận diện và bỏ qua chính xác các file không bị thay đổi trong repository.


2. Cách tiếp cận và Lý do lựa chọn#

Nhằm loại bỏ yếu tố sai số do thao tác thủ công, đồng thời đảm bảo quy trình kiểm thử mang tính minh bạch, lặp lại được nhiều lần (reproducible), nhóm đã phát triển một công cụ script tự động hóa hoàn chỉnh: scripts/demo_idempotent_replay.py.

Kịch bản tự động hóa kiểm chứng:#

  1. Lấy mẫu trạng thái hệ thống trước Replay: thiết lập kết nối trực tiếp đến MongoDB và Neo4j để truy vấn và lưu giữ tổng số lượng document tổng thể cũng như trạng thái metadata hiện tại của file mục tiêu.

  2. Chọn file demo có CPG nodes thực sự: script ưu tiên chọn các file lớn có ClassDef/FunctionDef thực sự (ở đây là modeling_bert.py) để đảm bảo Neo4j có số lượng node CPG lớn, việc kiểm chứng mới có ý nghĩa thực sự.

  3. Tiêm sự kiện thay đổi có kiểm soát: script tiến hành tự động chèn một dòng comment chứa mốc timestamp vào đuôi file mục tiêu để tạo sự khác biệt ở hàm băm (hash).

  4. Kích hoạt quy trình Incremental Parsing: thực thi lệnh python parser/main.py parse-changed --output-target kafka để ép buộc hệ thống chỉ tập trung quét và đẩy lên Kafka duy nhất event thuộc về file vừa bị chỉnh sửa.

  5. Kiểm tra trạng thái đối chiếu sau Replay với Assertion cứng:

    • MongoDB: truy xuất lại document của tệp mục tiêu để xác minh rằng các tham số siêu dữ liệu (content_hash, line_count) đã được cập nhật bản mới. Quan trọng nhất, kiểm tra hàm count_documents() phải trả về con số bất biến là 4.611 (minh chứng cho việc cơ chế Upsert hoạt động hoàn hảo).

    • Neo4j: Thực hiện assert nodes_after == nodes_before - nếu số node thay đổi, script báo lỗi ngay lập tức. Kiểm chứng thành công khi ASSERTION PASSED xuất hiện cùng số node cụ thể.

  6. Cơ chế Rollback (Khôi phục nguyên trạng): script đảm nhiệm việc xóa nội dung vừa chèn, trả tệp mục tiêu về tình trạng nguyên gốc ban đầu và tự động đồng bộ lại state store.


3. Mã nguồn & Kết quả thực thi#

Lệnh gọi kịch bản thực nghiệm tự động qua môi trường CLI:

python scripts/demo_idempotent_replay.py

Chuỗi log quá trình xử lý của kịch bản trên Terminal:

======================================================================
BẮT ĐẦU KỊCH BẢN DEMO IDEMPOTENT REPLAY
======================================================================
[INFO] File demo được chọn: modeling_bert.py (57,364 bytes)
[INFO] Target File được chọn để sửa đổi demo: transformers/src/transformers/models/bert/modeling_bert.py
[INFO] [MongoDB Trước Replay] Tổng số documents: 4611
       -> Doc 'transformers/src/transformers/models/bert/modeling_bert.py': line_count=1395, hash=edcdb558ad...
[INFO] [Neo4j Trước Replay] Số node của file 'transformers/src/transformers/models/bert/modeling_bert.py': 1729

[INFO] Đang thực hiện sửa đổi file (thêm comment tag replay)...
[INFO] Đang thực thi: python parser/main.py parse-changed --output-target kafka
Parsed 1 file(s), skipped 0, errors 0, wrote 1729 node event(s) and 4859 edge event(s).

[INFO] Chờ 4 giây để Kafka Connect (Neo4j Sink) và Spark Streaming (MongoDB) xử lý event...

======================================================================
KẾT QUẢ KIỂM CHỨNG SAU REPLAY (VERIFICATION RESULT)
======================================================================
[INFO] [MongoDB Sau Replay] Tổng số documents: 4611
       -> Doc 'transformers/src/transformers/models/bert/modeling_bert.py': line_count=1396, hash=f61ea80781...
[SUCCESS] XÁC MINH MONGODB: Document đã được UPSERT thành công, KHÔNG TẠO TRÙNG RECORD!
[INFO] [Neo4j Sau Replay] Số node của file 'transformers/src/transformers/models/bert/modeling_bert.py': 1729
[SUCCESS] XÁC MINH NEO4J (ASSERTION PASSED): nodes_before=1729 == nodes_after=1729. Cypher MERGE không gây nhân bản node!
[INFO] Xác nhận Neo4j: 1729 node CPG đã được MERGE an toàn, không có node rác nào được tạo ra.

[INFO] Đang khôi phục file về nội dung ban đầu...
[SUCCESS] Khôi phục hoàn tất! Kịch bản demo kết thúc thành công.

Lý do chọn modeling_bert.py thay vì __init__.py: File này chứa gần 1.400 dòng code với cấu trúc phức tạp, cung cấp số lượng node CPG lớn (1729 nodes) trước khi replay. Nhờ vậy, câu lệnh assert nodes_after == nodes_before mới thực sự kiểm chứng được rằng Cypher MERGE có khả năng hợp nhất trên quy mô lớn mà không nhân bản node.


Bảng Tổng Hợp Before → After#

Đây là bằng chứng quan trọng nhất của Task 6 - chứng minh toàn bộ cơ chế Idempotent hoạt động đúng:

Chỉ số

Trước Replay

Sau Replay

Kết luận

Neo4j Nodes

1729

1729

Giữ nguyên

Neo4j Edges

4859

4859

Giữ nguyên

MongoDB Documents

4.611

4.611

Không tăng

Modified File

Version 1 (1395 dòng)

Version 2 (1396 dòng)

Cập nhật

Giải thích kết quả:

  • Neo4j Nodes/Edges không thay đổi: Cypher MERGE chỉ cập nhật thuộc tính (SET) trên node/edge đã tồn tại, không tạo bản sao mới.

  • MongoDB Documents không tăng: Upsert với _id = file_path ghi đè document cũ, không tạo record mới.

  • Modified File cập nhật: content_hashline_count mới trong MongoDB Compass xác nhận file được xử lý lại thành công.

Kiểm chứng cơ chế Spark Checkpoint skip offset (theo yêu cầu §1.7): Sau khi replay, hệ thống ghi nhận rằng Spark Streaming job chỉ xử lý duy nhất 1 micro-batch mới từ topic cpg_metadata_events (chứa event của file đã thay đổi). Toàn bộ 4.610 offset còn lại đã được commit vào checkpoint từ lần chạy đầu tiên, do đó Spark hoàn toàn bỏ qua (skip) mà không tái xử lý.


4. Minh chứng thực tế#

Minh chứng 1: Trạng thái các dịch vụ Docker#

Trạng thái Docker Container

Hình 6.1: Các container Kafka, Neo4j, MongoDB và Kafka Connect đều đang ở trạng thái Healthy/Running.

Minh chứng 2: Kết quả chạy Script Idempotent Replay trên Terminal#

Output Terminal Script Replay

Hình 6.2: Màn hình console minh chứng luồng chạy hoàn hảo của Script Idempotent Replay, báo cáo Assertion Passed cho cả MongoDB và Neo4j.

Minh chứng 3: MongoDB & Neo4j TRƯỚC Replay#

MongoDB Document Trước Replay

Hình 6.3a: Document của modeling_bert.py trong MongoDB trước replay ghi nhận 1395 dòng (line_count).

MongoDB Tổng Document Trước Replay

Hình 6.3b: Tổng số document trong MongoDB trước replay là 4.611.

Neo4j Node Trước Replay

Hình 6.3c: Neo4j ghi nhận số node ban đầu của tệp mục tiêu là 1729 nodes.

Minh chứng 4: MongoDB & Neo4j SAU Replay#

MongoDB Document Sau Replay

Hình 6.4a: Document trong MongoDB được Upsert cập nhật content_hash và line_count mới (1396 dòng) nhưng giữ nguyên trường _id.

MongoDB Tổng Document Sau Replay

Hình 6.4b: Tổng số document trong MongoDB sau replay giữ nguyên bất biến là 4.611.

Neo4j Node Sau Replay

Hình 6.4c: Neo4j xác thực nodes_before == nodes_after (1729 nodes), không sinh ra node rác.

Minh chứng 5: Thư mục Spark Checkpoint#

Spark Checkpoint Directory

Hình 6.5: Cấu trúc thư mục Spark Checkpoint ghi nhận offset đã xử lý.


5. Tự nhận xét và Đúc kết#

Khép lại kịch bản kiểm chứng tự động hóa này, nhóm rất hài lòng khi thấy toàn bộ các mảnh ghép kiến trúc của hệ thống - từ Parser cục bộ, xương sống Kafka, kết nối Neo4j Connector đến luồng Spark Streaming - đã liên kết và phối hợp nhịp nhàng đúng như tính toán lý thuyết.

Thử thách kỹ thuật quan trọng khi thiết kế script demo nằm ở logic bảo vệ kho mã nguồn: sau khi chèn dữ liệu kiểm thử, kịch bản phải có cơ chế khôi phục lại (rollback) chính xác mã nguồn gốc để tránh làm thay đổi (pollute) Git repository. Giải pháp của nhóm là bọc khối logic thực thi thao tác file vào cấu trúc quản lý lỗi try...finally. Với cách tiếp cận này, ngay cả khi người dùng ngắt tiến trình đột ngột (Ctrl+C), mã nguồn mục tiêu vẫn luôn được bảo đảm khôi phục về trạng thái nguyên bản an toàn.


6. Hướng dẫn chạy lại#

  1. Kích hoạt toàn bộ hạ tầng Docker Compose (bao gồm Kafka, Neo4j, MongoDB) và tiến trình Spark Streaming ở chế độ nền:

    docker compose up -d
    
  2. Yêu cầu hệ thống khởi chạy script kiểm chứng tự động:

    python scripts/demo_idempotent_replay.py
    
  3. Theo dõi bảng kết quả kiểm chứng in ra trên màn hình Terminal và đối chiếu với dữ liệu Database.