用Kafka + ClickHouse构建实时数据分析平台:独立开发者的“数据引擎“实战

用Kafka + ClickHouse构建实时数据分析平台:独立开发者的"数据引擎"实战

为什么独立开发者需要实时数据分析

早期产品用"Google Analytics + 手工SQL查询"就够了。但当你的产品有"10万+用户"或"需要实时推荐/告警"时,你需要"专业的数据分析架构"。

典型场景:

  1. 实时仪表盘:显示"当前在线用户数"、"今日新增注册"、"实时收入"
  2. 用户行为分析:追踪"用户在哪个步骤掉落"、"哪些功能最常用"
  3. 告警系统:当"错误率突然升高"或"支付成功率下降"时立即通知

架构选型:为什么是Kafka + ClickHouse?

Kafka(消息队列):解耦"事件采集"和"事件处理"

如果你的产品每秒产生1000个事件(用户点击、API调用、错误日志),直接写入数据库会"压垮数据库"。

Kafka的作用:

  1. 缓冲:事件先写入Kafka,再由消费者慢慢处理
  2. 持久化:Kafka把事件保存7-30天,即使消费者挂了,数据也不丢失
  3. 可重放:可以"重新消费"历史事件(如"重新计算昨天的统计数据")

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):

  1. Upstash Kafka(Serverless Kafka,按写入量计费,免费额度高)
  2. ClickHouse Cloud(有免费额度,或自托管在VPS)

监控:

  • Kafka:Kafka UI(开源Web UI)
  • ClickHouse:Grafana + ClickHouse datasource

结论:Kafka + ClickHouse是"数据驱动产品"的基石

你不需要在"第一天"就上这套架构。但当你的产品有"10万+用户"或"需要实时数据分析"时,Kafka + ClickHouse是"性价比最高"的选择。

最小可行架构(MVP):

  1. 用Upstash Kafka(无需自己维护)
  2. 用ClickHouse Cloud(免费额度)
  3. 用Supabase pgvector做"用户行为相似度搜索"(可选)

这套架构,能支撑"每秒1万事件"的写入量——对独立开发者来说,足够用到"月收入$100K"的阶段了。

下一步: 把你的"用户行为数据"从Google Analytics迁移到"自建分析平台"——这样你就能"拥有数据",且能做"GA做不到的分析"(如"实时推荐"、"自定义告警")。

© 版权声明

相关文章