Báo cáo & Mã nguồn Đồ án 04 môn Nhập môn Dữ liệu Lớn (CSC14118) - Trường Đại học Khoa học Tự nhiên, ĐHQG-HCM.
Đồ án xây dựng đường ống dữ liệu (pipeline) xử lý luồng theo mô hình Event-Driven Architecture nhằm trích xuất Đồ thị thuộc tính mã nguồn (Code Property Graph - CPG) từ kho mã nguồn Python huggingface/transformers (~4.611 file Python).
Dữ liệu CPG (gồm AST, CFG, DFG) và siêu dữ liệu (metadata) được phân tải real-time qua Apache Kafka, sau đó được tiêu thụ và nạp song song vào 2 hệ thống lưu trữ:
MERGE lũy đẳng._id = file_path và quản lý Checkpoint an toàn.📖 Báo cáo chi tiết trực tuyến (Jupyter Book): https://htvien.github.io/bigdata-lab04/
Nhóm 4 - BigBang
| Họ và tên | MSSV | Nhiệm vụ Kỹ thuật | Phụ trách Báo cáo |
|---|---|---|---|
| Võ Ngọc Tiến | 23120370 | Shallow clone repo, viết CPG Parser (AST/CFG/DFG) & thuật toán Deterministic Stable ID | Chương 1 & 2: Repository Discovery & CPG Parser Service |
| Võ Thành Minh Tuệ | 23120398 | Thiết kế 4 Kafka topic, lập trình Producer, cấu hình Neo4j Kafka Connector (Direct Sink) | Chương 3 & 4: Kafka Topic Design & Neo4j Ingestion |
| Đỗ Phước Vinh | 23120405 | Phát triển Spark Structured Streaming tiêu thụ metadata, thiết kế schema MongoDB, quản lý Checkpoint | Chương 5: Spark Streaming & MongoDB Ingestion |
| Huỳnh Trọng Viên | 23120403 | Triển khai Docker Compose, viết script tự động kiểm thử Idempotent Replay, xuất bản Jupyter Book | Chương 6 & Hệ thống: Architecture, Idempotent Replay & Integrated Demo |
flowchart TD
subgraph Input["1. Input Source Code"]
HF["Hugging Face Transformers Repo<br/>(shallow-clone .py files)"]
end
subgraph ParserService["2. Parser Service (Python)"]
FileDisc["File Discovery & Exclusion Filter"]
ASTParser["CPG & AST Parser<br/>(AST / CFG / DFG / Call Edges)"]
StateStore["Parser State Store<br/>(content_hash & mtime)"]
EventWriter["Kafka Event Writer<br/>(confluent-kafka + LZ4)"]
HF --> FileDisc
FileDisc --> ASTParser
ASTParser <--> StateStore
ASTParser --> EventWriter
end
subgraph Broker["3. Apache Kafka Cluster (KRaft Mode)"]
T_Node["Topic: cpg_node_events"]
T_Edge["Topic: cpg_edge_events"]
T_Meta["Topic: cpg_metadata_events"]
T_Err["Topic: Error_parser (DLQ)"]
EventWriter --> T_Node
EventWriter --> T_Edge
EventWriter --> T_Meta
EventWriter --> T_Err
end
subgraph StorageSink["4. Databases & Streaming Sinks"]
subgraph Neo4jSink["Neo4j Kafka Connector (Direct Sink)"]
SinkConn["Kafka Connect Worker<br/>(neo4j-kafka-connect)"]
Neo4jDB[("Neo4j Graph Database<br/>(CPG Topology / Cypher MERGE)")]
T_Node --> SinkConn
T_Edge --> SinkConn
SinkConn -->|Direct Ingestion without Spark| Neo4jDB
end
subgraph SparkSink["Spark Structured Streaming Sink"]
SparkJob["Spark Streaming Job<br/>(foreachBatch & Checkpoint)"]
MongoDB[("MongoDB Database<br/>(Source Code Metadata / Upsert _id)")]
T_Meta --> SparkJob
SparkJob -->|_id = file_path| MongoDB
SparkJob -.->|Malformed JSON| T_Err
end
end
bigdata-lab04/
├── book/ # Nguồn tài liệu Jupyter Book (Markdown + Config)
│ ├── _config.yml # Cấu hình Jupyter Book
│ ├── _toc.yml # Mục lục các chương
│ ├── architecture.md # Tổng quan kiến trúc
│ ├── chapter1_discovery.md # Task 1: Repository Discovery
│ ├── chapter2_cpg_parser.md # Task 2: Incremental CPG Parser
│ ├── chapter3_kafka_design.md # Task 3: Kafka Topic Design
│ ├── chapter4_neo4j_ingestion.md # Task 4: Neo4j Ingestion
│ ├── chapter5_spark_mongodb.md # Task 5: Spark Streaming & MongoDB
│ ├── chapter6_idempotent_replay.md # Task 6: Idempotent Replay
│ └── acknowledgments.md # Lời cảm ơn thầy cô
├── docs/ # Bản sao đồng bộ của book/ phục vụ GitHub Pages
├── kafka/ # Cấu hình Kafka & Neo4j Sink Connector
│ ├── neo4j-sink-config.json # Cấu hình Cypher MERGE cho Neo4j Sink
│ └── producer.py # Kafka Producer module
├── parser/ # Parser Service trích xuất CPG
│ ├── cpg_parser.py # Phân tích CPG: AST, CFG, DFG, Call edges
│ ├── file_discovery.py # Quét & lọc tệp Python (.py)
│ ├── event_writer.py # Ghi events ra file JSONL hoặc Kafka
│ ├── models.py # Schema dataclass cho NodeEvent, EdgeEvent, MetadataEvent
│ ├── id_utils.py # Hàm stable_node_id / stable_edge_id deterministic
│ ├── config.py # Cấu hình topic, path, schema version
│ ├── main.py # CLI chính của Parser Service
│ └── state_store.py # Quản lý content_hash cho incremental parse
├── spark/ # Spark Structured Streaming Job
│ └── metadata_streaming.py # Spark job đọc metadata từ Kafka -> MongoDB
├── scripts/ # Kịch bản tự động hóa & kiểm thử
│ ├── demo_idempotent_replay.py # Script kiểm thử Idempotent Replay tự động
│ └── generate_svg_arch.py # Script sinh sơ đồ kiến trúc SVG
├── docker-compose.yml # Cấu hình hạ tầng (Kafka, Neo4j, MongoDB, Connect)
├── requirements.txt # Danh sách thư viện Python phụ thuộc
└── README.md # Tài liệu hướng dẫn sử dụng dự án
# Clone repository của nhóm
git clone https://github.com/htvien/bigdata-lab04.git
cd bigdata-lab04
# Clone shallow copy repository Hugging Face Transformers (input)
git clone --depth 1 https://github.com/huggingface/transformers.git transformers
Trên Windows (PowerShell):
python -m venv venv
.\venv\Scripts\activate
pip install -r requirements.txt
Trên Linux / macOS:
python3 -m venv venv
source venv/bin/activate
pip install -r requirements.txt
docker compose up -d
Sau khi khởi động, kiểm tra trạng thái các container:
docker compose ps
Các cổng dịch vụ mặc định:
localhost:9092http://localhost:8080http://localhost:8083http://localhost:7474 (User: neo4j / Password: password)localhost:7687localhost:27017http://localhost:4040 (khi Spark job đang chạy)python parser/main.py discover
Kết quả: Sinh file output/discovered_python_files.json chứa thông tin 4.611 tệp .py.
python parser/main.py parse-all --output-target kafka
Kết quả: Phân tích toàn bộ mã nguồn, sinh ra 2.029.594 node events và 6.062.224 edge events đẩy vào 4 topic Kafka:
cpg_node_eventscpg_edge_eventscpg_metadata_eventsError_parser(Lưu ý: Quá trình khởi tạo connector nạp trực tiếp Node/Edge vào Neo4j đã được tự động hóa khi chạy Docker Compose. Bạn không cần chạy lệnh POST để đăng ký nữa).
Kiểm tra trạng thái Connector:
curl -s http://localhost:8083/connectors/neo4j-sink-connector/status
python spark/metadata_streaming.py
Kết quả: Spark tiêu thụ luồng metadata từ Kafka, thực hiện Upsert vào MongoDB collection bigdata_lab04.metadata với khóa chính _id = file_path và ghi vết checkpoint tại output/checkpoints/metadata.
Chạy script kiểm thử tự động kịch bản thay đổi 1 tệp và replay:
python scripts/demo_idempotent_replay.py
Kết quả xác minh:
content_hash mới nhưng tổng số document giữ nguyên 4.611 (Upsert thành công, không trùng record).MERGE, không gây nhân bản node.Báo cáo đồ án được xuất bản dưới dạng Jupyter Book tương tác tại địa chỉ:
👉 https://htvien.github.io/bigdata-lab04/
Nội dung báo cáo gồm 7 phần:
Nhóm BigBang xin trân trọng gửi lời cảm ơn chân thành đến TS. Nguyễn Ngọc Thảo (giảng dạy lý thuyết và hướng dẫn đồ án), cùng các thầy trợ giảng ThS. Huỳnh Lâm Hải Đăng và thầy Trần Huy Bân (hướng dẫn thực hành) đã tận tình giúp đỡ nhóm trong suốt môn học Nhập môn Dữ liệu Lớn (CSC14118).
Môn Nhập môn Dữ liệu Lớn (CSC14118) - Trường Đại học Khoa học tự nhiên, ĐHQG-HCM - 2026
Python
99.9%
Báo cáo & Mã nguồn Đồ án 04 môn Nhập môn Dữ liệu Lớn (CSC14118) - Trường Đại học Khoa học Tự nhiên, ĐHQG-HCM.
Đồ án xây dựng đường ống dữ liệu (pipeline) xử lý luồng theo mô hình Event-Driven Architecture nhằm trích xuất Đồ thị thuộc tính mã nguồn (Code Property Graph - CPG) từ kho mã nguồn Python huggingface/transformers (~4.611 file Python).
Dữ liệu CPG (gồm AST, CFG, DFG) và siêu dữ liệu (metadata) được phân tải real-time qua Apache Kafka, sau đó được tiêu thụ và nạp song song vào 2 hệ thống lưu trữ:
MERGE lũy đẳng._id = file_path và quản lý Checkpoint an toàn.📖 Báo cáo chi tiết trực tuyến (Jupyter Book): https://htvien.github.io/bigdata-lab04/
Nhóm 4 - BigBang
| Họ và tên | MSSV | Nhiệm vụ Kỹ thuật | Phụ trách Báo cáo |
|---|---|---|---|
| Võ Ngọc Tiến | 23120370 | Shallow clone repo, viết CPG Parser (AST/CFG/DFG) & thuật toán Deterministic Stable ID | Chương 1 & 2: Repository Discovery & CPG Parser Service |
| Võ Thành Minh Tuệ | 23120398 | Thiết kế 4 Kafka topic, lập trình Producer, cấu hình Neo4j Kafka Connector (Direct Sink) | Chương 3 & 4: Kafka Topic Design & Neo4j Ingestion |
| Đỗ Phước Vinh | 23120405 | Phát triển Spark Structured Streaming tiêu thụ metadata, thiết kế schema MongoDB, quản lý Checkpoint | Chương 5: Spark Streaming & MongoDB Ingestion |
| Huỳnh Trọng Viên | 23120403 | Triển khai Docker Compose, viết script tự động kiểm thử Idempotent Replay, xuất bản Jupyter Book | Chương 6 & Hệ thống: Architecture, Idempotent Replay & Integrated Demo |
flowchart TD
subgraph Input["1. Input Source Code"]
HF["Hugging Face Transformers Repo<br/>(shallow-clone .py files)"]
end
subgraph ParserService["2. Parser Service (Python)"]
FileDisc["File Discovery & Exclusion Filter"]
ASTParser["CPG & AST Parser<br/>(AST / CFG / DFG / Call Edges)"]
StateStore["Parser State Store<br/>(content_hash & mtime)"]
EventWriter["Kafka Event Writer<br/>(confluent-kafka + LZ4)"]
HF --> FileDisc
FileDisc --> ASTParser
ASTParser <--> StateStore
ASTParser --> EventWriter
end
subgraph Broker["3. Apache Kafka Cluster (KRaft Mode)"]
T_Node["Topic: cpg_node_events"]
T_Edge["Topic: cpg_edge_events"]
T_Meta["Topic: cpg_metadata_events"]
T_Err["Topic: Error_parser (DLQ)"]
EventWriter --> T_Node
EventWriter --> T_Edge
EventWriter --> T_Meta
EventWriter --> T_Err
end
subgraph StorageSink["4. Databases & Streaming Sinks"]
subgraph Neo4jSink["Neo4j Kafka Connector (Direct Sink)"]
SinkConn["Kafka Connect Worker<br/>(neo4j-kafka-connect)"]
Neo4jDB[("Neo4j Graph Database<br/>(CPG Topology / Cypher MERGE)")]
T_Node --> SinkConn
T_Edge --> SinkConn
SinkConn -->|Direct Ingestion without Spark| Neo4jDB
end
subgraph SparkSink["Spark Structured Streaming Sink"]
SparkJob["Spark Streaming Job<br/>(foreachBatch & Checkpoint)"]
MongoDB[("MongoDB Database<br/>(Source Code Metadata / Upsert _id)")]
T_Meta --> SparkJob
SparkJob -->|_id = file_path| MongoDB
SparkJob -.->|Malformed JSON| T_Err
end
end
bigdata-lab04/
├── book/ # Nguồn tài liệu Jupyter Book (Markdown + Config)
│ ├── _config.yml # Cấu hình Jupyter Book
│ ├── _toc.yml # Mục lục các chương
│ ├── architecture.md # Tổng quan kiến trúc
│ ├── chapter1_discovery.md # Task 1: Repository Discovery
│ ├── chapter2_cpg_parser.md # Task 2: Incremental CPG Parser
│ ├── chapter3_kafka_design.md # Task 3: Kafka Topic Design
│ ├── chapter4_neo4j_ingestion.md # Task 4: Neo4j Ingestion
│ ├── chapter5_spark_mongodb.md # Task 5: Spark Streaming & MongoDB
│ ├── chapter6_idempotent_replay.md # Task 6: Idempotent Replay
│ └── acknowledgments.md # Lời cảm ơn thầy cô
├── docs/ # Bản sao đồng bộ của book/ phục vụ GitHub Pages
├── kafka/ # Cấu hình Kafka & Neo4j Sink Connector
│ ├── neo4j-sink-config.json # Cấu hình Cypher MERGE cho Neo4j Sink
│ └── producer.py # Kafka Producer module
├── parser/ # Parser Service trích xuất CPG
│ ├── cpg_parser.py # Phân tích CPG: AST, CFG, DFG, Call edges
│ ├── file_discovery.py # Quét & lọc tệp Python (.py)
│ ├── event_writer.py # Ghi events ra file JSONL hoặc Kafka
│ ├── models.py # Schema dataclass cho NodeEvent, EdgeEvent, MetadataEvent
│ ├── id_utils.py # Hàm stable_node_id / stable_edge_id deterministic
│ ├── config.py # Cấu hình topic, path, schema version
│ ├── main.py # CLI chính của Parser Service
│ └── state_store.py # Quản lý content_hash cho incremental parse
├── spark/ # Spark Structured Streaming Job
│ └── metadata_streaming.py # Spark job đọc metadata từ Kafka -> MongoDB
├── scripts/ # Kịch bản tự động hóa & kiểm thử
│ ├── demo_idempotent_replay.py # Script kiểm thử Idempotent Replay tự động
│ └── generate_svg_arch.py # Script sinh sơ đồ kiến trúc SVG
├── docker-compose.yml # Cấu hình hạ tầng (Kafka, Neo4j, MongoDB, Connect)
├── requirements.txt # Danh sách thư viện Python phụ thuộc
└── README.md # Tài liệu hướng dẫn sử dụng dự án
# Clone repository của nhóm
git clone https://github.com/htvien/bigdata-lab04.git
cd bigdata-lab04
# Clone shallow copy repository Hugging Face Transformers (input)
git clone --depth 1 https://github.com/huggingface/transformers.git transformers
Trên Windows (PowerShell):
python -m venv venv
.\venv\Scripts\activate
pip install -r requirements.txt
Trên Linux / macOS:
python3 -m venv venv
source venv/bin/activate
pip install -r requirements.txt
docker compose up -d
Sau khi khởi động, kiểm tra trạng thái các container:
docker compose ps
Các cổng dịch vụ mặc định:
localhost:9092http://localhost:8080http://localhost:8083http://localhost:7474 (User: neo4j / Password: password)localhost:7687localhost:27017http://localhost:4040 (khi Spark job đang chạy)python parser/main.py discover
Kết quả: Sinh file output/discovered_python_files.json chứa thông tin 4.611 tệp .py.
python parser/main.py parse-all --output-target kafka
Kết quả: Phân tích toàn bộ mã nguồn, sinh ra 2.029.594 node events và 6.062.224 edge events đẩy vào 4 topic Kafka:
cpg_node_eventscpg_edge_eventscpg_metadata_eventsError_parser(Lưu ý: Quá trình khởi tạo connector nạp trực tiếp Node/Edge vào Neo4j đã được tự động hóa khi chạy Docker Compose. Bạn không cần chạy lệnh POST để đăng ký nữa).
Kiểm tra trạng thái Connector:
curl -s http://localhost:8083/connectors/neo4j-sink-connector/status
python spark/metadata_streaming.py
Kết quả: Spark tiêu thụ luồng metadata từ Kafka, thực hiện Upsert vào MongoDB collection bigdata_lab04.metadata với khóa chính _id = file_path và ghi vết checkpoint tại output/checkpoints/metadata.
Chạy script kiểm thử tự động kịch bản thay đổi 1 tệp và replay:
python scripts/demo_idempotent_replay.py
Kết quả xác minh:
content_hash mới nhưng tổng số document giữ nguyên 4.611 (Upsert thành công, không trùng record).MERGE, không gây nhân bản node.Báo cáo đồ án được xuất bản dưới dạng Jupyter Book tương tác tại địa chỉ:
👉 https://htvien.github.io/bigdata-lab04/
Nội dung báo cáo gồm 7 phần:
Nhóm BigBang xin trân trọng gửi lời cảm ơn chân thành đến TS. Nguyễn Ngọc Thảo (giảng dạy lý thuyết và hướng dẫn đồ án), cùng các thầy trợ giảng ThS. Huỳnh Lâm Hải Đăng và thầy Trần Huy Bân (hướng dẫn thực hành) đã tận tình giúp đỡ nhóm trong suốt môn học Nhập môn Dữ liệu Lớn (CSC14118).
Môn Nhập môn Dữ liệu Lớn (CSC14118) - Trường Đại học Khoa học tự nhiên, ĐHQG-HCM - 2026
Python
99.9%