Đồ án 04 - Spark Streaming#
Báo cáo đồ án 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.
Nhóm BigBang xây dựng đường ống dữ liệu (pipeline) xử lý luồng 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, sau đó nạp kết quả vào hai hệ thống lưu trữ độc lập là Neo4j và MongoDB thông qua Apache Kafka và Spark Structured Streaming.
Thông tin môn học#
Môn học |
Nhập môn Dữ liệu lớn (CSC14118) |
Giảng viên |
TS. Nguyễn Ngọc Thảo, TS. Lê Ngọc Thành |
Trợ giảng |
ThS. Huỳnh Lâm Hải Đăng, thầy Trần Huy Bân |
Nhóm |
BigBang |
Repository |
|
Phân công thành viên nhóm#
Họ và tên |
MSSV |
Nhiệm vụ |
Phụ trách viết báo cáo |
|---|---|---|---|
Võ Ngọc Tiến |
23120370 |
Cấu hình Shallow clone repository, viết CPG Parser (AST/CFG/DFG) & thiết kế 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ý cơ chế checkpoint. |
Chương 5: Spark Streaming & MongoDB Ingestion |
Huỳnh Trọng Viên |
23120403 |
Triển khai hạ tầng Docker Compose, viết script kiểm thử tự động Idempotent Replay, cấu hình Jupyter Book & xuất bản Pages. |
Chương 6 & Hệ thống: Architecture, Idempotent Replay & Tổng hợp |
Cấu trúc báo cáo#
Báo cáo được chia thành 7 phần chính, trình bày theo đúng tiến trình nhóm đã triển khai thực tế trên hệ thống:
Kiến trúc tổng quan: Sơ đồ khối hệ thống và diễn giải chi tiết luồng luân chuyển dữ liệu.
Chương 1: Repository Cloning & File Discovery: Chiến lược shallow clone và thuật toán duyệt, lọc mã nguồn Python.
Chương 2: Incremental CPG Parser Service: Kỹ thuật bóc tách AST, CFG, DFG và thuật toán định danh ID cố định bằng mã băm SHA-256.
Chương 3: Kafka Topic Design & Schema Events: Kiến trúc 4 topic Kafka chuyên biệt và quy chuẩn schema định dạng tin nhắn.
Chương 4: Neo4j Ingestion: Cơ chế nạp trực tiếp dữ liệu đồ thị từ Kafka vào Neo4j thông qua truy vấn Cypher MERGE, không qua Spark.
Chương 5: Spark Streaming & MongoDB Ingestion: Xử lý luồng siêu dữ liệu qua Spark và thực hiện Upsert vào MongoDB, kết hợp quản lý Checkpoint.
Chương 6: Idempotent Replay Verification: Thực nghiệm kiểm chứng kịch bản cập nhật dữ liệu lũy đẳng, đảm bảo không sinh ra bản ghi trùng lặp.