Chương 5: Spark Structured Streaming & MongoDB Ingestion#
Thành viên phụ trách: Đỗ Phước Vinh
1. Bối cảnh và Mục tiêu#
Trọng tâm của chương này là thiết kế và triển khai một tiến trình Spark Structured Streaming với nhiệm vụ tiêu thụ liên tục luồng sự kiện siêu dữ liệu (cpg_metadata_events) từ hàng đợi Kafka. Khối dữ liệu sau khi được phân tích sẽ được nạp vào kho lưu trữ NoSQL MongoDB (bigdata_lab04.metadata).
Điểm then chốt trong yêu cầu kỹ thuật là job Spark phải duy trì được cơ chế Checkpoint (lưu vết trạng thái) nghiêm ngặt. Khả năng này đảm bảo hệ thống sở hữu năng lực chịu lỗi (fault-tolerance): nếu tiến trình bị sập hay ngắt điện, lúc tái khởi động, Spark có thể dò tìm lại chính xác chỉ mục (offset) đã đọc cuối cùng trong Kafka để tiếp tục công việc, triệt tiêu rủi ro bỏ sót tin nhắn hoặc xử lý lặp dữ liệu.
2. Cách tiếp cận và Lý do lựa chọn#
Để đáp ứng song song tiêu chí về hiệu năng luồng streaming và khả năng chịu lỗi bền bỉ, nhóm đã thiết lập cấu trúc xử lý như sau:
Ứng dụng toán tử Upsert trong MongoDB thông qua khóa chính
_id = file_path: Kỹ thuật chèn (insert) thông thường sẽ liên tục tạo ra các document rác khi mã nguồn bị phân tích nhiều lần. Nhóm khắc phục bằng cách ánh xạ chủ động đường dẫn của tệp (file_path) thành thuộc tính khóa chính_idcủa MongoDB. Thông qua việc thiết lập MongoDB Spark Connector sang modeupdate(cụ thể:spark.mongodb.write.operationType = updatevàspark.mongodb.write.idFieldList = _id), hệ thống duy trì được tính tỷ lệ 1-1 chặt chẽ: mỗi file Python vật lý chỉ tương ứng với duy nhất một document trong bộ sưu tập (collection). Các thay đổi siêu dữ liệu mới sẽ được MongoDB tự động ghi đè lên document hiện hữu.Cơ chế lưu vết Checkpoint tập trung: Nhóm khai thác tính năng quản lý state tự động của Spark Structured Streaming. Spark sẽ định kỳ ghi log về các offset của Kafka cùng trạng thái hoàn thành của từng micro-batch vào cấu trúc thư mục vật lý
output/checkpoints/metadata. Nhờ cơ sở dữ liệu trạng thái này, mọi vòng đời gián đoạn của quá trình Spark đều có thể được khôi phục nguyên trạng một cách tin cậy.
3. Mã nguồn & Kết quả thực thi#
Cấu trúc logic xử lý micro-batch và khai báo vị trí checkpoint trong spark/metadata_streaming.py:
from pyspark.sql.functions import col, from_json
# ... (Khởi tạo SparkSession và Schema) ...
# BƯỚC 1: Đọc luồng dữ liệu từ Kafka
df = (
spark.readStream.format("kafka")
.option("kafka.bootstrap.servers", args.kafka_broker)
.option("subscribe", KAFKA_TOPIC)
.option("startingOffsets", "earliest")
.option("failOnDataLoss", "false")
.load()
)
# BƯỚC 2: Phân tích JSON từ cột value
parsed_df = (
df.selectExpr("CAST(value AS STRING)")
.select(from_json(col("value"), schema).alias("data"))
.select("data.*")
)
# BƯỚC 3: Map `file_path` thành `_id` (Cơ sở cho Idempotent Upsert)
upsert_df = parsed_df.withColumn("_id", col("file_path"))
# BƯỚC 4: Ghi luồng vào MongoDB và bật Checkpoint
query = (
upsert_df.writeStream.format("mongodb")
.option("spark.mongodb.connection.uri", args.mongo_uri)
.option("spark.mongodb.database", args.mongo_db)
.option("spark.mongodb.collection", MONGO_COLLECTION)
.option("spark.mongodb.write.operationType", "update")
.option("spark.mongodb.write.idFieldList", "_id")
.option("checkpointLocation", CHECKPOINT_DIR)
.outputMode("append")
.start()
)
Script đo lường số lượng Document thực tế trong MongoDB:
python -c "from pymongo import MongoClient; client = MongoClient('mongodb://127.0.0.1:27017'); print('MongoDB metadata count:', client['bigdata_lab04']['metadata'].count_documents({}))"
Output thực tế (Terminal):
MongoDB metadata count: 4611
Thống kê xác thực toàn bộ 4.611 tập tin metadata đã hoàn tất chu trình được Spark tiêu thụ từ hệ thống trung gian Kafka và Upsert thành công vào cấu trúc lưu trữ của MongoDB.
4. Minh chứng thực tế#
Minh chứng 1: Bảng điều khiển Spark Structured Streaming UI (http://localhost:4040)#

Hình 5.1: Màn hình giao diện quản trị Spark UI (tab Structured Streaming) phản ánh trạng thái “RUNNING” của streaming query, đang trực tiếp tiêu thụ và xử lý các micro-batch.
Minh chứng 2: Trạng thái Document trong MongoDB TRƯỚC khi chạy Replay#

Hình 5.2: Giao diện MongoDB Compass hiển thị chi tiết một document trước khi cập nhật. Lưu ý các giá trị ban đầu: content_hash là hash_lan_1_abc123 và line_count là 450.
Minh chứng 3: Trạng thái Document trong MongoDB SAU khi chạy Replay (Minh chứng Upsert)#

Hình 5.3: Document đã được cơ chế Upsert ghi đè thông tin mới thành công (content_hash đổi thành hash_lan_2_xyz890, line_count thành 460) mà vẫn giữ nguyên khóa chính _id, không sinh thêm bản ghi rác.
Minh chứng 4: Cấu trúc phòng vệ của thư mục Spark Checkpoint (output/checkpoints/metadata)#

Hình 5.4: Hệ thống phân bổ trạng thái lưu vết hoàn chỉnh của Spark. Các thư mục con offsets/ và commits/ lưu trữ liên tục các tệp tin trạng thái (từ 0 đến 5…) giúp đảm bảo khả năng chịu lỗi phục hồi luồng.
5. Tự nhận xét và Đúc kết#
Về những mặt đã hoạt động tốt, cơ chế Idempotent Upsert với khóa chính _id kết hợp với cấu hình operationType = update của MongoDB Spark Connector đã hoạt động trơn tru ngay từ lần chạy đầu tiên. Bằng chứng rõ ràng nhất ở Hình 5.2 và 5.3 đã cho thấy hệ thống ngăn chặn thành công việc tạo bản ghi trùng lặp khi chạy kịch bản Incremental Replay.
Việc triển khai Spark cục bộ trên môi trường hệ điều hành Windows thường đi kèm với những đặc thù lỗi hạ tầng phức tạp, và dự án lần này cũng không ngoại lệ.
Trở ngại kỹ thuật đầu tiên liên quan đến lỗi môi trường Java HADOOP_HOME and hadoop.home.dir are unset. Do hệ thống Windows thiếu vắng các driver hệ thống tập tin phân tán (HDFS) tự nhiên, nhóm đã nghiên cứu giải pháp bổ sung bộ nhị phân giả lập winutils.exe và hadoop.dll, khởi tạo cấu trúc hadoop_home/bin và khai báo biến môi trường vào Python script để đáp ứng yêu cầu bộ phụ thuộc Hadoop của Spark.
Vấn đề kỹ thuật tiếp theo diễn ra khi vô tình kích hoạt đồng thời 2 terminal cùng chạy tiến trình python spark/metadata_streaming.py. Engine của Spark đã phát ra ngoại lệ khóa luồng Multiple streaming queries are concurrently using offset checkpoint. Đây là cơ chế bảo vệ tính toàn vẹn (lock mechanism) của Spark để cấm các job khác nhau ghi đè lên cùng một thư mục checkpoint. Nhóm đã xử lý bằng cách kết thúc tiến trình dư thừa và Spark tiếp tục vận hành ổn định.
6. Hướng dẫn chạy lại#
Thẩm định trạng thái hoạt động của các container nền tảng (MongoDB và Kafka):
docker compose ps mongodb kafka
Khởi tạo tác vụ Spark Streaming xử lý dữ liệu liên tục:
python spark/metadata_streaming.py
Truy cập để quản lý và theo dõi thông lượng xử lý của Spark UI qua địa chỉ:
http://localhost:4040