Chấm dứt nỗi lo hỗn loạn Log: Xây dựng Pipeline chuẩn hóa với Pydantic

Programming tutorial - IT technology blog
Programming tutorial - IT technology blog

Cuộc gọi đánh thức lúc 2:15 sáng

Vào lúc 2:15 sáng thứ Ba, điện thoại của tôi không chỉ rung; nó gào thét. PagerDuty báo cáo có sự gia tăng đột biến 450 lỗi mỗi phút trong dịch vụ thanh toán (checkout service). Tôi lảo đảo bước đến bàn làm việc, mở trình tổng hợp log và nhìn chằm chằm vào một bức tường nhiễu kỹ thuật số. Đó là một đống hỗn độn hoàn toàn.

Dashboard lúc đó như một “nghĩa địa” của những dòng văn bản không thể đọc nổi. Thay vì dữ liệu sạch và có thể tìm kiếm, tôi thấy ba dịch vụ khác nhau đang “hét” lên bằng ba ngôn ngữ khác nhau:

2023-10-27 02:14:58 INFO [auth_service] User ID: 4502 - login success
{"level": "error", "timestamp": "2023-10-27T02:15:01Z", "message": "Database timeout", "service": "checkout"}
[CRITICAL] 02:15:05 - payment_gateway - Connection refused - IP: 10.0.0.5

Việc tìm kiếm của chúng tôi thất bại vì một nửa số log không ở định dạng JSON. Tôi đã lãng phí 45 phút vật lộn với các mẫu Regex mong manh chỉ để tách biệt một địa chỉ IP bị lỗi. Đây chính là cái giá tiềm ẩn của log không cấu trúc. Khi cơ sở hạ tầng của bạn gặp sự cố, bạn không nên phải ngồi viết các bộ parser (phân tích cú pháp) ngay tức thì.

Tại sao Log lại bị phân mảnh?

Sự phân mảnh log thường không phải là kết quả của kỹ thuật kém. Đó là một tác dụng phụ tự nhiên của việc mở rộng kiến trúc microservices. Các nhóm khác nhau sử dụng các công cụ khác nhau. Một nhóm thích Loguru, nhóm khác trung thành với module logging tiêu chuẩn, và dịch vụ Java cũ kỹ ở trong góc thì chỉ sử dụng System.out.println().

Đến khi những log này đi tới Elasticsearch hoặc Loki, hệ thống đã bị quá tải. Việc bỏ qua bước chuẩn hóa tại điểm nạp dữ liệu (ingestion point) sẽ biến hệ thống giám sát của bạn thành một gánh nặng. Dashboard bị hỏng. Cảnh báo bỏ lỡ các đợt tăng đột biến quan trọng. Các công cụ phân tích tự động đơn giản là bị “nghẹn” trước dữ liệu không nhất quán. Về cơ bản, bạn đang lái một chiếc máy bay với buồng lái đầy những đồng hồ đo bị vỡ vụn.

Đánh giá giải pháp: Regex vs. Parsing thủ công vs. Pydantic

Tôi đã đánh giá ba cách để làm sạch dữ liệu này. Mỗi cách đều có những ưu và nhược điểm riêng.

1. Sử dụng Regex

Bạn có thể viết một script Python đồ sộ chứa đầy các biểu thức chính quy (regular expressions) để bắt mọi mẫu log. Mặc dù nhanh, nhưng Regex là một cơn ác mộng về bảo trì. Nếu một lập trình viên thêm dù chỉ một khoảng trắng vào thông điệp log, toàn bộ pipeline sẽ bị hỏng. Nó rất dễ gãy, khó đọc và khó kiểm thử.

2. Phân tích Dictionary thủ công

Sử dụng string.split() và ánh xạ dictionary thủ công có thể hoạt động cho các trường hợp đơn giản, nhưng nó không cung cấp khả năng xác thực. Nếu một “User ID” đến dưới dạng chuỗi thay vì số nguyên, các công cụ phân tích hạ nguồn của bạn sẽ bị sập vài giờ sau đó. Bạn chỉ đang đẩy vấn đề đi xa hơn chứ không giải quyết triệt để.

3. Sử dụng Pydantic Model

Pydantic là một thư viện xác thực tận dụng Python type hints. Nó không chỉ phân tích dữ liệu; nó thực thi một schema nghiêm ngặt. Nếu dữ liệu không khớp, nó sẽ cho bạn biết chính xác lý do tại sao. Nó tự động xử lý ép kiểu (type casting)—chuyển đổi chuỗi “200” thành số nguyên 200—và xuất mọi thứ sang định dạng JSON tiêu chuẩn.

Xây dựng Pipeline chuẩn hóa

Tôi đã xây dựng một lớp chuẩn hóa trung tâm sử dụng Pydantic để giải quyết vấn đề này. Mục tiêu là tiếp nhận bất kỳ chuỗi thô hoặc dictionary hỗn độn nào và ép nó vào một đối tượng NormalizedLog có kiểu dữ liệu nghiêm ngặt. Cách tiếp cận này đã biến quy trình debug của chúng tôi từ một trò chơi đoán mò thành một hoạt động chính xác.

Bước 1: Định nghĩa Base Schema

Chúng ta cần một cấu trúc chuẩn mà mọi log phải tuân theo. Điều này đảm bảo tính nhất quán trên toàn bộ hệ thống.

from pydantic import BaseModel, Field, field_validator
from datetime import datetime, timezone
from typing import Optional, Any
import uuid

class NormalizedLog(BaseModel):
    log_id: str = Field(default_factory=lambda: str(uuid.uuid4()))
    timestamp: datetime
    level: str
    service_name: str
    message: str
    payload: Optional[dict[str, Any]] = None

    @field_validator('level')
    @classmethod
    def normalize_level(cls, v: str) -> str:
        # Chuẩn hóa level thành chữ hoa và loại bỏ khoảng trắng
        return v.upper().strip()

Bước 2: Tạo các Parser chuyên dụng

Tiếp theo, chúng ta tạo logic để xử lý các định dạng khác nhau. Tôi sử dụng chiến lược dự phòng (fallback): thử phân tích log dưới dạng JSON trước, sau đó chuyển sang Regex cho các log kiểu cũ (legacy).

import json
import re

class LogNormalizer:
    # Regex cho định dạng: [CRITICAL] 02:15:05 - payment_gateway - Message
    LEGACY_PATTERN = re.compile(r"\[(?P<level>\w+)\] (?P<time>[\d:]+) - (?P<service>[\w_]+) - (?P<msg>.*)")

    def normalize(self, raw_data: str) -> NormalizedLog:
        try:
            data = json.loads(raw_data)
            return NormalizedLog(
                timestamp=data.get("timestamp", datetime.now(timezone.utc)),
                level=data.get("level", "INFO"),
                service_name=data.get("service", "không_xác_định"),
                message=data.get("message", ""),
                payload=data
            )
        except json.JSONDecodeError:
            pass

        match = self.LEGACY_PATTERN.search(raw_data)
        if match:
            groups = match.groupdict()
            return NormalizedLog(
                timestamp=datetime.now(timezone.utc),
                level=groups['level'],
                service_name=groups['service'],
                message=groups['msg']
            )
        
        return NormalizedLog(
            timestamp=datetime.now(timezone.utc),
            level="CHƯA_XÁC_ĐỊNH",
            service_name="chưa_phân_tích",
            message=raw_data
        )

Bước 3: Xử lý hiệu suất cao

Trong môi trường production, bạn có thể phải xử lý 10.000 log mỗi giây. Pydantic v2 đóng vai trò quan trọng ở đây vì logic cốt lõi của nó được viết bằng Rust, giúp nó nhanh hơn tới 20 lần so với v1. Đối với khối lượng nạp dữ liệu lớn, tôi bao bọc logic này trong một async worker để ngăn chặn tình trạng thắt nút cổ chai.

import asyncio

async def process_logs(raw_logs: list[str]):
    normalizer = LogNormalizer()
    normalized_data = []
    
    for raw in raw_logs:
        # Chuyển đổi sang đối tượng Pydantic và sau đó sang chuỗi JSON
        entry = normalizer.normalize(raw)
        normalized_data.append(entry.model_dump_json())
    
    await save_to_storage(normalized_data)

async def save_to_storage(data):
    # Tải lên hàng loạt (batch upload) tới Elasticsearch, Loki, hoặc S3
    print(f"Đang lưu trữ {len(data)} bản ghi đã chuẩn hóa.")

Kết quả thực tế

Điều làm cho hệ thống này trở nên mạnh mẽ là khả năng xử lý lỗi một cách êm đẹp. Nếu một log hoàn toàn không thể nhận dạng, nó vẫn được bao bọc trong một đối tượng NormalizedLog với level là “CHƯA_XÁC_ĐỊNH”. Pipeline nạp dữ liệu của bạn không bao giờ bị sập. Bạn chỉ cần thiết lập cảnh báo cho service_name == "chưa_phân_tích" để phát hiện và khắc phục các định dạng log mới.

Đến khi sự cố tiếp theo xảy ra, dashboard của chúng tôi đã hoạt động cực kỳ chính xác. Chúng tôi lọc theo service_name, sắp xếp theo timestamp và đi sâu vào payload mà không cần đụng đến Regex. Chúng tôi đã xác định được nguyên nhân gốc rễ—cạn kiệt pool kết nối cơ sở dữ liệu—trong vòng chưa đầy ba phút. Nếu bạn vẫn đang vật lộn với văn bản thô, đã đến lúc xây dựng một schema. Hãy tặng cho bản thân tương lai của bạn món quà là những giấc ngủ ngon.

Share: