
Xây dựng tác nhân AI cục bộ truyền phát
Khi nói về các tác nhân AI (AI agents), thuật ngữ "streaming" được sử dụng theo hai cách khác nhau. Cần làm rõ sự hiểu biết về vấn đề này.
Xây dựng tác nhân AI cục bộ truyền phát - KDnuggets
"Truyền phát" (streaming) được sử dụng theo hai cách khác nhau khi mọi người nói về các tác nhân AI. Hầu hết các hướng dẫn chỉ xây dựng một trong số đó. Đôi khi, nó có nghĩa là tác nhân tiêu thụ một luồng sự kiện trực tiếp thay vì chờ đợi ai đó nhập tin nhắn. Đôi khi, nó có nghĩa là đầu ra của chính tác nhân được truyền phát từng token một thay vì xuất hiện tất cả cùng lúc sau một khoảng dừng dài. Bản dựng này thực hiện cả hai điều đó, một cách có chủ đích, vì chúng giải quyết hai vấn đề khác nhau, và một tác nhân luôn hoạt động thực sự hữu ích cần cả hai vấn đề này được giải quyết.
Khung khái niệm đáng mượn ở đây đến từ cái thường được gọi là tác nhân môi trường (ambient agent), một loại tác nhân mà LangChain mô tả là được kích hoạt bởi các sự kiện thay vì bởi một tin nhắn từ con người, và Bộ công cụ phát triển tác nhân của Google (Google's Agent Development Kit) mô tả từ phía hạ tầng theo cùng một cách: các tác nhân được đánh thức bởi một thứ gì đó đến trên một luồng, không phải ngồi sau một cuộc gọi yêu cầu-phản hồi. Kịch bản cho bản dựng này là cụ thể và hoàn toàn có thật: một tác nhân cục bộ theo dõi nguồn cấp dữ liệu chỉnh sửa trực tiếp, công khai của Wikipedia, không yêu cầu khóa API, và đưa ra lý do về những chỉnh sửa nào trông giống như phá hoại, chạy hoàn toàn trên máy tính của bạn thông qua Ollama. Mọi dòng mã dưới đây đều được viết, sau đó thực sự được kiểm tra, trước khi đưa vào bài viết này.
Đây là các điều kiện tiên quyết của bạn:
Python 3.11 trở lên
Ollama được cài đặt cục bộ, với một mô hình đã được tải (ollama pull llama3.1:8b, hoặc bất kỳ mô hình nào hỗ trợ đầu ra JSON có cấu trúc)
pip install fastapi uvicorn httpx pydantic ollama sse-starlette
Không có khóa API, không có tài khoản đám mây và không có chi phí nào ngoài tiền điện của bạn. Kết nối mạng đi duy nhất mà dịch vụ này thực hiện là đến điểm cuối EventStreams công khai của Wikipedia, không yêu cầu xác thực.
# Quyết định thiết kế quan trọng nhất
Luồng chỉnh sửa của Wikipedia không phải là một dòng chảy nhỏ giọt. Vào một ngày hoạt động, nó đẩy vài chỉnh sửa mỗi giây trên tất cả các phiên bản ngôn ngữ kết hợp. Nếu đưa từng chỉnh sửa đó cho một mô hình ngôn ngữ, hai điều sẽ xảy ra cùng lúc: bạn sẽ tiêu tốn tài nguyên tính toán của máy cho những chỉnh sửa không bao giờ thú vị ngay từ đầu, và tác nhân sẽ bị tụt lại phía sau luồng trực tiếp mà nó phải theo dõi, điều này làm mất đi toàn bộ mục đích của việc xây dựng một thứ "luôn hoạt động".
Giải pháp là một phễu hai giai đoạn, và đó là ý tưởng quan trọng nhất trong bản dựng này:
Giai đoạn một là phép toán Python đơn giản, chi phí thấp, chạy trên mọi sự kiện mà không cần mô hình nào cả: chỉnh sửa này đã xóa bao nhiêu byte, người dùng này đã thực hiện bao nhiêu chỉnh sửa trong vài phút qua? Phần lớn các chỉnh sửa đều nhàm chán, và việc phát hiện sự nhàm chán là miễn phí.
Giai đoạn hai, LLM cục bộ thực tế, chỉ hoạt động đối với một phần nhỏ các sự kiện vượt ngưỡng ở giai đoạn một. Đây là nguyên tắc tương tự đằng sau bất kỳ hệ thống giám sát tốt nào: các bộ lọc rẻ tiền ở phía trước, suy luận phức tạp được dành riêng cho các ứng viên vượt qua vòng loại.
// Cấu trúc thư mục
streaming-local-agent/
├── src/
│ ├── __init__.py
│ ├── config.py
│ ├── schemas.py
│ ├── stream_source.py
│ ├── filters.py
│ ├── agent.py
│ ├── broadcaster.py
│ └── main.py
├── tests/
│ └── test_filters.py
├── requirements.txt
└── .env.example
Mỗi tệp tương ứng chính xác với một giai đoạn của quy trình được mô tả ở trên, điều này giúp toàn bộ hệ thống dễ hiểu và dễ kiểm tra độc lập, đúng như cách nó được xây dựng cho bài viết này.
# Phần Xây dựng 1: Trình tiêu thụ luồng sự kiện
Dịch vụ EventStreams của Wikipedia đẩy các chỉnh sửa dưới dạng Server-Sent Events qua HTTP thuần túy. Không có khóa, không có bắt tay ngoài một yêu cầu GET thông thường vẫn mở.
# src/stream_source.py
import asyncio
import json
import re
import time
from typing import AsyncIterator, Optional
import httpx
from .schemas import RecentChangeEvent
from . import config
# Wikipedia không gửi cờ "người dùng này ẩn danh" rõ ràng trên luồng này;
# các chỉnh sửa ẩn danh được gán cho địa chỉ IP của người chỉnh sửa thay vì tên người dùng,
# vì vậy, một tên người dùng có dạng IP là cách bạn phát hiện ra một người dùng ẩn danh trên thực tế.
_IPV4_RE = re.compile(r"^\d{1,3}(\.\d{1,3}){3}$")
_IPV6_RE = re.compile(r"^[0-9A-Fa-f:]+:[0-9A-Fa-f:]+$")
def is_anonymous_user(username: str) -> bool:
return bool(_IPV4_RE.match(username) or _IPV6_RE.match(username))
def parse_sse_line(line: str) -> Optional[dict]:
"""Các khung SSE dữ liệu dưới dạng các dòng có tiền tố 'data: '. Các dòng bình luận
(bắt đầu bằng ':') và các dòng giữ kết nối trống rỗng thường xuất hiện trên luồng này
và nên được bỏ qua một cách lặng lẽ, không được coi là lỗi."""
if not line or line.startswith(":"):
return None
if line.startswith("data:"):
raw = line[len("data:"):].strip()
if not raw:
return None
try:
return json.loads(raw)
except json.JSONDecodeError:
return None
return None
def to_event(raw: dict) -> Optional[RecentChangeEvent]:
"""Chuyển đổi tải trọng Wikimedia thô thành lược đồ chuẩn hóa của chúng tôi.
Trả về None cho các loại sự kiện mà chúng tôi không quan tâm thay vì
gây ra lỗi, vì một luồng có khối lượng lớn như vậy liên tục bao gồm các định dạng
mà chúng tôi không theo dõi."""
if raw.get("type") != "edit":
return None
length = raw.get("length") or {}
if "old" not in length or "new" not in length:
return None
return RecentChangeEvent(
wiki=raw.get("wiki", "unknown"),
user=raw.get("user", "unknown"),
title=raw.get("title", "unknown"),
Nguồn tin: KDnuggets — Tác giả: Shittu Olumide. Bản dịch tiếng Việt do AI thực hiện, có thể có sai sót.