Kiến trúc tổng thể hệ thống#
Chương này trình bày sơ đồ kiến trúc và chi tiết luồng luân chuyển dữ liệu của hệ thống CPG Streaming Pipeline do nhóm thiết kế và xây dựng.
1. Sơ đồ kiến trúc luồng dữ liệu#
Hệ thống được thiết kế theo hướng Event-Driven Architecture (Kiến trúc hướng sự kiện), phân chia thành 4 cụm xử lý độc lập để tối ưu hóa hiệu năng và khả năng mở rộng:
Hình 1: Sơ đồ kiến trúc tổng thể 4 tầng (Input Source Code, Parser Service, Apache Kafka Cluster, Neo4j & MongoDB Streaming Sinks).
2. Nguyên lý hoạt động của từng thành phần#
Input & Parser Service:
Nhóm sử dụng kỹ thuật shallow-clone để tải kho mã nguồn
huggingface/transformersvề môi trường cục bộ làm dữ liệu đầu vào.Parser Service (phát triển bằng Python) chịu trách nhiệm duyệt tuần tự qua từng file
.py, phân tích cấu trúc để trích xuất Cây cú pháp trừu tượng (AST), Luồng điều khiển (CFG) và Luồng dữ liệu (DFG).Điểm cốt lõi là mỗi phần tử (node/edge) khi được trích xuất sẽ được gán một định danh cố định (Deterministic Stable ID) dựa trên thuật toán băm SHA-256. Điều này đảm bảo tính nhất quán tuyệt đối của ID khi phân tích lại cùng một tệp.
Apache Kafka Message Broker:
Đóng vai trò là trục xương sống trung chuyển sự kiện, vận hành ở chế độ KRaft (loại bỏ sự phụ thuộc vào Zookeeper).
Broker nhận luồng dữ liệu từ Parser Service và phân phối vào 4 topic riêng biệt:
cpg_node_events,cpg_edge_events,cpg_metadata_eventsvàError_parserđể cách ly các luồng xử lý.
Neo4j Kafka Connect Direct Sink:
Tiêu thụ dữ liệu trực tiếp từ 2 topic cấu trúc (
cpg_node_eventsvàcpg_edge_events).Luồng sự kiện được nạp thẳng vào cơ sở dữ liệu đồ thị Neo4j mà không cần đi qua tầng Spark. Giải pháp này giúp giảm thiểu đáng kể độ trễ và độ phức tạp của hạ tầng.
Kết hợp sử dụng câu lệnh Cypher
MERGEđể ngăn chặn triệt để hiện tượng nhân bản dữ liệu (data duplication) khi pipeline được thực thi lại.
Spark Structured Streaming & MongoDB Sink:
Mạch xử lý này đảm nhận việc tiêu thụ luồng sự kiện từ topic
cpg_metadata_events.Dữ liệu siêu cấu trúc được Spark xử lý và ghi vào cơ sở dữ liệu MongoDB (
bigdata_lab04.metadata) thông qua cơ chế Upsert, sử dụng đường dẫn file (file_path) làm khóa chính_id.Trạng thái của luồng streaming được bảo toàn thông qua cơ chế Checkpoint, lưu trữ an toàn tại thư mục
output/checkpoints/metadata.