用Kafka + ClickHouse构建实时数据分析平台:独立开发者的“数据引擎“实战
用Kafka + ClickHouse构建实时数据分析平台:独立开发者的"数据引擎"实战
为什么独立开发者需要实时数据分析
早期产品用"Google Analytics + 手工SQL查询"就够了。但当你的产品有"10万+用户"或"需要实时推荐/告警"时,你需要"专业的数据分析架构"。
典型场景:
- 实时仪表盘:显示"当前在线用户数"、"今日新增注册"、"实时收入"
- 用户行为分析:追踪"用户在哪个步骤掉落"、"哪些功能最常用"
- 告警系统:当"错误率突然升高"或"支付成功率下降"时立即通知
架构选型:为什么是Kafka + ClickHouse?
Kafka(消息队列):解耦"事件采集"和"事件处理"
如果你的产品每秒产生1000个事件(用户点击、API调用、错误日志),直接写入数据库会"压垮数据库"。
Kafka的作用:
- 缓冲:事件先写入Kafka,再由消费者慢慢处理
- 持久化:Kafka把事件保存7-30天,即使消费者挂了,数据也不丢失
- 可重放:可以"重新消费"历史事件(如"重新计算昨天的统计数据")
ClickHouse(列式数据库):实时分析查询
ClickHouse是"为分析查询设计的数据库"——它用"列式存储"和"向量化执行",让"聚合查询"(如"计算今日UV"、"统计每个功能的使用次数")比PostgreSQL快10-100倍。
对比:PostgreSQL vs ClickHouse(1000万行数据)
| 查询类型 | PostgreSQL | ClickHouse | 差异 |
|---|---|---|---|
| COUNT(*) | 2.5s | 0.05s | ClickHouse快50倍 |
| GROUP BY + 聚合 | 8.2s | 0.3s | ClickHouse快27倍 |
| 按时间范围查询 | 1.8s | 0.08s | ClickHouse快22倍 |
实战:用Kafka + ClickHouse构建"用户行为分析平台"
第一步:用Docker Compose搭建本地开发环境
# docker-compose.yml
version: '3.8'
services:
zookeeper:
image: confluentinc/cp-zookeeper:7.5.0
hostname: zookeeper
container_name: zookeeper
ports:
- "2181:2181"
environment:
ZOOKEEPER_CLIENT_PORT: 2181
ZOOKEEPER_TICK_TIME: 2000
kafka:
image: confluentinc/cp-kafka:7.5.0
hostname: kafka
container_name: kafka
depends_on:
- zookeeper
ports:
- "29092:29092"
- "9092:9092"
environment:
KAFKA_BROKER_ID: 1
KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181'
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
clickhouse:
image: clickhouse/clickhouse-server:latest
hostname: clickhouse
container_name: clickhouse
ports:
- "8123:8123" # HTTP API
- "9000:9000" # Native TCP
volumes:
- ./clickhouse-data:/var/lib/clickhouse
启动:
docker-compose up -d
第二步:创建Kafka Topic(事件流)
# 进入Kafka容器
docker exec -it kafka kafka-topics --create \
--topic user-events \
--bootstrap-server kafka:29092 \
--partitions 3 \
--replication-factor 1
第三步:在ClickHouse创建事件表
-- 连接到ClickHouse(用clickhouse-client或HTTP API)
CREATE TABLE user_events
(
event_time DateTime64(3),
user_id String,
event_type String, -- 'page_view', 'button_click', 'signup', etc.
page_url String,
metadata String, -- JSON字符串,存储额外信息
session_id String
)
ENGINE = MergeTree()
PARTITION BY toYYYYMMDD(event_time)
ORDER BY (user_id, event_time);
-- 创建物化视图(实时聚合)
CREATE MATERIALIZED VIEW daily_active_users
ENGINE = SummingMergeTree()
PARTITION BY toYYYYMMDD(event_date)
ORDER BY (event_date)
AS SELECT
toDate(event_time) as event_date,
uniq(user_id) as dau
FROM user_events
WHERE event_type = 'page_view'
GROUP BY event_date;
关键: ClickHouse的MergeTree引擎会自动"按ORDER BY排序数据"——这让"时间范围查询"极快。
事件采集:在前端/后端发送事件到Kafka
后端采集(Node.js + kafkajs):
// lib/kafkaProducer.ts
import { Kafka, Producer } from 'kafkajs';
const kafka = new Kafka({
clientId: 'my-product',
brokers: ['localhost:9092'],
});
let producer: Producer | null = null;
export async function getProducer() {
if (!producer) {
producer = kafka.producer();
await producer.connect();
}
return producer;
}
export async function sendEvent(event: {
userId: string;
eventType: string;
pageUrl?: string;
metadata?: Record<string, any>;
sessionId: string;
}) {
const producer = await getProducer();
await producer.send({
topic: 'user-events',
messages: [
{
key: event.userId, // 相同用户的事件发到同一个partition(保证顺序)
value: JSON.stringify({
event_time: new Date().toISOString(),
user_id: event.userId,
event_type: event.eventType,
page_url: event.pageUrl || '',
metadata: JSON.stringify(event.metadata || {}),
session_id: event.sessionId,
}),
},
],
});
}
在API Route里调用:
// app/api/signup/route.ts
export async function POST(req: Request) {
const { email, password } = await req.json();
const user = await db.users.create({ email, password });
// 发送注册事件
await sendEvent({
userId: user.id,
eventType: 'signup',
metadata: { method: 'email' },
sessionId: req.headers.get('X-Session-ID') || 'unknown',
});
return Response.json(user);
}
前端采集(JavaScript):
// lib/analytics.ts
export function trackEvent(eventType: string, metadata?: Record<string, any>) {
// 发送到你的后端API(后端再写入Kafka)
fetch('/api/events', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
event_type: eventType,
page_url: window.location.pathname,
metadata,
session_id: getSessionId(), // 从cookie或localStorage获取
}),
}).catch(err => console.error('Failed to track event:', err));
}
// 自动追踪页面浏览
if (typeof window !== 'undefined') {
trackEvent('page_view');
// 监听路由变化(如果是SPA)
window.addEventListener('popstate', () => trackEvent('page_view'));
}
Kafka消费者:把事件写入ClickHouse
// consumers/clickhouseWriter.ts
import { Kafka, Consumer } from 'kafkajs';
import { createClient } from '@clickhouse/client';
const kafka = new Kafka({ brokers: ['localhost:9092'] });
const consumer: Consumer = kafka.consumer({ groupId: 'clickhouse-writer' });
const clickhouse = createClient({
url: 'http://localhost:8123',
database: 'default',
});
async function runConsumer() {
await consumer.connect();
await consumer.subscribe({ topic: 'user-events', fromBeginning: false });
await consumer.run({
eachMessage: async ({ topic, partition, message }) => {
const event = JSON.parse(message.value!.toString());
// 写入ClickHouse
await clickhouse.insert({
table: 'user_events',
values: [{
event_time: event.event_time,
user_id: event.user_id,
event_type: event.event_type,
page_url: event.page_url,
metadata: event.metadata,
session_id: event.session_id,
}],
format: 'JSONEachRow',
});
console.log(`Inserted event: ${event.event_type} for user ${event.user_id}`);
},
});
}
runConsumer().catch(console.error);
优化:批量写入(提升性能)
// 改进版:批量写入(每100ms或1000条事件写入一次)
class BatchWriter {
private batch: any[] = [];
private timer: NodeJS.Timeout | null = null;
async add(event: any) {
this.batch.push(event);
if (!this.timer) {
this.timer = setTimeout(() => this.flush(), 100);
}
if (this.batch.length >= 1000) {
await this.flush();
}
}
private async flush() {
if (this.batch.length === 0) return;
const toInsert = [...this.batch];
this.batch = [];
this.timer = null;
await clickhouse.insert({
table: 'user_events',
values: toInsert,
format: 'JSONEachRow',
});
console.log(`Flushed ${toInsert.length} events to ClickHouse`);
}
}
实时查询:构建API给Dashboard用
// app/api/analytics/dau/route.ts
import { createClient } from '@clickhouse/client';
const clickhouse = createClient({ url: 'http://localhost:8123' });
export async function GET(req: Request) {
const { searchParams } = new URL(req.url);
const days = parseInt(searchParams.get('days') || '7');
// 查询每日活跃用户(DAU)
const result = await clickhouse.query({
query: `
SELECT
event_date,
dau
FROM daily_active_users
WHERE event_date >= today() - INTERVAL ${days} DAY
ORDER BY event_date DESC
`,
format: 'JSONEachRow',
});
const data = await result.json();
return Response.json(data);
}
在Dashboard展示(React + Recharts):
// components/DAUChart.tsx
import { useEffect, useState } from 'react';
import { LineChart, Line, XAxis, YAxis, CartesianGrid, Tooltip } from 'recharts';
export function DAUChart() {
const [data, setData] = useState<any[]>([]);
useEffect(() => {
fetch('/api/analytics/dau?days=30')
.then(res => res.json())
.then(setData);
}, []);
return (
<LineChart width={800} height={400} data={data}>
<XAxis dataKey="event_date" />
<YAxis />
<CartesianGrid strokeDasharray="3 3" />
<Tooltip />
<Line type="monotone" dataKey="dau" stroke="#8884d8" />
</LineChart>
);
}
成本与运维:独立开发者的可行方案
成本分析(云端部署):
| 组件 | 自托管(VPS) | 托管服务(Confluent Cloud + ClickHouse Cloud) |
|---|---|---|
| Kafka | 免费(自己维护) | $0.10/GB 写入 + $0.05/GB 存储 |
| ClickHouse | 免费(自己维护) | $0.10/GB 存储 + $0.05/GB 查询 |
| 运维成本 | 高(需要懂Kafka/ClickHouse) | 低(托管服务) |
推荐方案(月收入<$10K):
- 用Upstash Kafka(Serverless Kafka,按写入量计费,免费额度高)
- 用ClickHouse Cloud(有免费额度,或自托管在VPS)
监控:
- Kafka:Kafka UI(开源Web UI)
- ClickHouse:Grafana + ClickHouse datasource
结论:Kafka + ClickHouse是"数据驱动产品"的基石
你不需要在"第一天"就上这套架构。但当你的产品有"10万+用户"或"需要实时数据分析"时,Kafka + ClickHouse是"性价比最高"的选择。
最小可行架构(MVP):
- 用Upstash Kafka(无需自己维护)
- 用ClickHouse Cloud(免费额度)
- 用Supabase pgvector做"用户行为相似度搜索"(可选)
这套架构,能支撑"每秒1万事件"的写入量——对独立开发者来说,足够用到"月收入$100K"的阶段了。
下一步: 把你的"用户行为数据"从Google Analytics迁移到"自建分析平台"——这样你就能"拥有数据",且能做"GA做不到的分析"(如"实时推荐"、"自定义告警")。