htvien/bigdata-lab04

0

stars

88

commits

Python

primary language

Jul 25, 2026

updated

README

Đồ án 04 - Incremental CPG Streaming Pipeline

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.

Jupyter Book Repository Python Kafka Spark Neo4j MongoDB


1. Giới thiệu tổng quan

Đồ á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ữ:

  • Neo4j Graph Database: Nạp trực tiếp qua Kafka Connect Neo4j Sink (không qua Spark) bằng các truy vấn Cypher MERGE lũy đẳng.
  • MongoDB NoSQL Database: Nạp qua Spark Structured Streaming với cơ chế Upsert khóa chính _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/


2. Phân công thành viên nhóm

Nhóm 4 - BigBang

Họ và tênMSSVNhiệm vụ Kỹ thuậtPhụ trách Báo cáo
Võ Ngọc Tiến23120370Shallow clone repo, viết CPG Parser (AST/CFG/DFG) & thuật toán Deterministic Stable IDChương 1 & 2: Repository Discovery & CPG Parser Service
Võ Thành Minh Tuệ23120398Thiế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 Vinh23120405Phát triển Spark Structured Streaming tiêu thụ metadata, thiết kế schema MongoDB, quản lý CheckpointChương 5: Spark Streaming & MongoDB Ingestion
Huỳnh Trọng Viên23120403Triển khai Docker Compose, viết script tự động kiểm thử Idempotent Replay, xuất bản Jupyter BookChương 6 & Hệ thống: Architecture, Idempotent Replay & Integrated Demo

3. Kiến trúc tổng thể hệ thống

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

4. Cấu trúc thư mục dự án

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

5. Yêu cầu môi trường & Cài đặt

Yêu cầu hệ thống

  • Hệ điều hành: Windows 10/11, macOS, hoặc Linux
  • Python: 3.11+
  • Docker Desktop: Đã bật Docker Engine
  • Java JDK: 11 hoặc 17 (cần thiết cho PySpark)

Bước 1: Clone dự án và repository đầu vào

# 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

Bước 2: Khởi tạo môi trường ảo Python

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

Bước 3: Khởi động các dịch vụ nền bằng Docker Compose

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:

  • Kafka Bootstrap Broker: localhost:9092
  • Kafka UI: http://localhost:8080
  • Kafka Connect REST API: http://localhost:8083
  • Neo4j Browser: http://localhost:7474 (User: neo4j / Password: password)
  • Neo4j Bolt Port: localhost:7687
  • MongoDB: localhost:27017
  • Spark Web UI: http://localhost:4040 (khi Spark job đang chạy)

6. Quy trình vận hành & Thực thi Pipeline

Bước 1: Khám phá danh sách tệp Python (Task 1)

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.

Bước 2: Phân tích CPG và đẩy luồng sự kiện lên Kafka (Task 2 & 3)

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_events
  • cpg_edge_events
  • cpg_metadata_events
  • Error_parser

Bước 3: Kiểm tra Neo4j Kafka Connector Sink (Task 4)

(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

Bước 4: Khởi chạy Spark Structured Streaming vào MongoDB (Task 5)

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.

Bước 5: Kiểm chứng Idempotent Replay (Task 6)

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:

  • Document MongoDB của tệp bị sửa đổi được cập nhật 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).
  • Neo4j được cập nhật qua Cypher MERGE, không gây nhân bản node.
  • Spark Checkpoint bỏ qua 4.610 offset cũ không bị thay đổi.

7. Báo cáo & Tài liệu trực tuyến

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:

  1. Kiến trúc tổng thể
  2. Chương 1: Repository Cloning & File Discovery
  3. Chương 2: Incremental CPG Parser Service
  4. Chương 3: Kafka Topic Design & Schema Events
  5. Chương 4: Neo4j Ingestion & Connector Configuration
  6. Chương 5: Spark Structured Streaming & MongoDB Ingestion
  7. Chương 6: Idempotent Replay Verification
  8. Lời Cảm Ơn

8. Lời cảm ơ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 Đăngthầ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

Contributors

htvien

56 commits

tue29092005a

16 commits

DPVinhIT

12 commits

tientien01

4 commits

htvien/bigdata-lab04

0

stars

88

commits

Python

primary language

Jul 25, 2026

updated

README

Đồ án 04 - Incremental CPG Streaming Pipeline

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.

Jupyter Book Repository Python Kafka Spark Neo4j MongoDB


1. Giới thiệu tổng quan

Đồ á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ữ:

  • Neo4j Graph Database: Nạp trực tiếp qua Kafka Connect Neo4j Sink (không qua Spark) bằng các truy vấn Cypher MERGE lũy đẳng.
  • MongoDB NoSQL Database: Nạp qua Spark Structured Streaming với cơ chế Upsert khóa chính _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/


2. Phân công thành viên nhóm

Nhóm 4 - BigBang

Họ và tênMSSVNhiệm vụ Kỹ thuậtPhụ trách Báo cáo
Võ Ngọc Tiến23120370Shallow clone repo, viết CPG Parser (AST/CFG/DFG) & thuật toán Deterministic Stable IDChương 1 & 2: Repository Discovery & CPG Parser Service
Võ Thành Minh Tuệ23120398Thiế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 Vinh23120405Phát triển Spark Structured Streaming tiêu thụ metadata, thiết kế schema MongoDB, quản lý CheckpointChương 5: Spark Streaming & MongoDB Ingestion
Huỳnh Trọng Viên23120403Triển khai Docker Compose, viết script tự động kiểm thử Idempotent Replay, xuất bản Jupyter BookChương 6 & Hệ thống: Architecture, Idempotent Replay & Integrated Demo

3. Kiến trúc tổng thể hệ thống

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

4. Cấu trúc thư mục dự án

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

5. Yêu cầu môi trường & Cài đặt

Yêu cầu hệ thống

  • Hệ điều hành: Windows 10/11, macOS, hoặc Linux
  • Python: 3.11+
  • Docker Desktop: Đã bật Docker Engine
  • Java JDK: 11 hoặc 17 (cần thiết cho PySpark)

Bước 1: Clone dự án và repository đầu vào

# 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

Bước 2: Khởi tạo môi trường ảo Python

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

Bước 3: Khởi động các dịch vụ nền bằng Docker Compose

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:

  • Kafka Bootstrap Broker: localhost:9092
  • Kafka UI: http://localhost:8080
  • Kafka Connect REST API: http://localhost:8083
  • Neo4j Browser: http://localhost:7474 (User: neo4j / Password: password)
  • Neo4j Bolt Port: localhost:7687
  • MongoDB: localhost:27017
  • Spark Web UI: http://localhost:4040 (khi Spark job đang chạy)

6. Quy trình vận hành & Thực thi Pipeline

Bước 1: Khám phá danh sách tệp Python (Task 1)

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.

Bước 2: Phân tích CPG và đẩy luồng sự kiện lên Kafka (Task 2 & 3)

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_events
  • cpg_edge_events
  • cpg_metadata_events
  • Error_parser

Bước 3: Kiểm tra Neo4j Kafka Connector Sink (Task 4)

(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

Bước 4: Khởi chạy Spark Structured Streaming vào MongoDB (Task 5)

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.

Bước 5: Kiểm chứng Idempotent Replay (Task 6)

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:

  • Document MongoDB của tệp bị sửa đổi được cập nhật 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).
  • Neo4j được cập nhật qua Cypher MERGE, không gây nhân bản node.
  • Spark Checkpoint bỏ qua 4.610 offset cũ không bị thay đổi.

7. Báo cáo & Tài liệu trực tuyến

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:

  1. Kiến trúc tổng thể
  2. Chương 1: Repository Cloning & File Discovery
  3. Chương 2: Incremental CPG Parser Service
  4. Chương 3: Kafka Topic Design & Schema Events
  5. Chương 4: Neo4j Ingestion & Connector Configuration
  6. Chương 5: Spark Structured Streaming & MongoDB Ingestion
  7. Chương 6: Idempotent Replay Verification
  8. Lời Cảm Ơn

8. Lời cảm ơ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 Đăngthầ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

Contributors

htvien

56 commits

tue29092005a

16 commits

DPVinhIT

12 commits

tientien01

4 commits

Languages

Python

99.9%