背景
大批量 SERP 数据需要实时处理 (监控 + 告警 + 数据湖),用 Kafka 作为 SERP 数据流的核心基础设施最自然。
# 启动 Kafka(本地)
docker run -d --name kafka -p 9092:9092 \
-e KAFKA_ADVERTISED_HOST=localhost \
apache/kafka
# 安装 Python client
pip install kafka-python requests
import json
import requests
from kafka import KafkaProducer
import time
producer = KafkaProducer(
bootstrap_servers="localhost:9092",
value_serializer=lambda v: json.dumps(v).encode("utf-8"),
)
def fetch_and_produce():
"""调 serpbase 抓 SERP → 推 Kafka"""
queries = ["python async", "rust tutorial", "react hooks"]
for q in queries:
r = requests.post(
"https://api.serpbase.dev/google/search",
headers={"X-API-Key": "sk_xxx"},
json={"q": q, "gl": "us", "num": 10},
timeout=10,
)
data = r.json()
for i, item in enumerate(data.get("organic", []), 1):
producer.send("serp-results", value={
"query": q,
"rank": i,
"title": item["title"],
"link": item["link"],
"ts": int(time.time()),
})
producer.flush()
fetch_and_produce()
import json
from kafka import KafkaConsumer
import requests
consumer = KafkaConsumer(
"serp-results",
bootstrap_servers="localhost:9092",
value_deserializer=lambda v: json.loads(v.decode("utf-8")),
)
def process_event(event):
"""每个 SERP 事件做实时处理"""
# 1. 实时分析(LLM 摘要)
if is_important_keyword(event["query"]):
summary = call_llm(event)
send_slack(f"SERP 变化:{summary}")
# 2. 写数据库
db.serp_results.insert(event)
# 3. 推下游(告警、数据湖)
if event["rank"] > 10:
alert_low_rank(event)
if event["rank"] < 3:
alert_top_rank(event)
# 推送到 Elasticsearch
es.index(index="serp-data", body=event)
for event in consumer:
process_event(event)
# 用 Avro / Protobuf 强 schema
{
"type": "record",
"name": "SerpResult",
"fields": [
{"name": "query", "type": "string"},
{"name": "rank", "type": "int"},
{"name": "title", "type": "string"},
{"name": "ts", "type": "long"},
]
}
producer = KafkaProducer(
transactional_id="serp-producer-1", # 启用事务
enable_idempotence=True,
)
consumer = KafkaConsumer(
"serp-results",
group_id="serp-processors", # 多个 consumer 自动分片
)
import fastavro
# 写
with open("schema.avsc") as f:
schema = json.load(f)
parsed = fastavro.parse_schema(schema)
buf = io.BytesIO()
fastavro.schemaless_writer(buf, parsed).write(record)
producer.send("serp-results", value=buf.getvalue())
// Java Kafka Streams 实时聚合
StreamsBuilder builder = new StreamsBuilder();
KStream<String, SerpResult> source = builder.stream("serp-results");
KTable<String, Long> rankByQuery = source
.groupBy((k, v) -> v.query)
.aggregate(
() -> 0L,
(k, v, acc) -> (long) v.rank,
Materialized.as("rank-by-query")
);
rankByQuery.toStream().to("rank-summary");
consumer = KafkaConsumer("serp-results", group_id="rank-monitor")
for event in consumer:
if event.rank > 10:
send_alert(event.query, event.rank)
# 实时推 Kafka → 5 分钟归档到 S3
from kafka import KafkaConsumer
consumer = KafkaConsumer("serp-results", group_id="archiver")
batch = []
for event in consumer:
batch.append(event)
if len(batch) >= 1000:
s3.put_object(Bucket="serp-data-lake", Key=f"raw/{int(time.time()) // 300 * 300}.json", Body=json.dumps(batch))
batch = []
KStream<String, SerpResult> events = builder.stream("serp-results");
KTable<Windowed<String>, Long> rankByHour = events
.groupBy((k, v) -> v.query)
.windowedBy(TimeWindows.of(Duration.ofHours(1)))
.aggregate(
() -> 0L,
(k, v, acc) -> Math.min(acc, v.rank),
Materialized.as("min-rank-by-hour")
);
from kafka import KafkaConsumer
import asyncio
async def process_event(event):
summary = await call_llm_async(event)
await send_slack_async(summary)
consumer = KafkaConsumer("serp-results", group_id="ai-processors")
for event in consumer:
asyncio.create_task(process_event(event))
{
"name": "serp-events-sink",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
"topics": "serp-results",
"connection.url": "jdbc:postgresql://localhost:5432/seo",
"auto.create": true,
}
}
| 指标 | 数值 |
|---|---|
| 每天调 serpbase | 10,000 次 |
| Kafka 消息吞吐 | 100K / 天 |
| 处理延迟 | < 100ms |
| 月成本 Kafka(本地) | $0(自部署) |
| 月成本 serpbase | $3.00 |
| 维度 | Kafka | RabbitMQ | Redis Streams |
|---|---|---|---|
| 吞吐 | 极高 | 高 | 中 |
| 持久化 | 强 | 中 | 弱 |
| 复杂 | 高 | 中 | 低 |
| 适合 | 大数据流 | 任务队列 | 简单流 |
serpbase + Kafka 流式管道:
Kafka 适合大数据 + 高吞吐,serpbase + Kafka 解决 1k-10k QPS 的实时处理需求。