这个项目能做什么
RealTime-ECommerce-Analytics(RTEA)是一个高性能的Python流式管道,从REST API捕获电商点击流数据,并通过Confluent Cloud Kafka路由到ClickHouse进行实时分析查询。
## 架构
该管道包含四个层级:
**数据源**——托管在Railway.app上的模拟点击流生成器API,每秒产生约10,000条JSON格式消息。
**消息中间件**——Confluent Cloud Kafka,Topic `clickstream-events` 含6个分区,采用SASL_SSL认证、Snappy压缩和7天保留策略。
**处理层**——Producer服务使用`kafka-python`并配置优化参数(16KB批次、gzip压缩、32MB缓冲区、限速10K消息/秒);Consumer服务采用批量处理(每批次1000条事件、自动提交偏移量、字段转换)。
**分析存储**——ClickHouse Cloud中基于MergeTree引擎的表,按`(timestamp, user_id)`排序,使用LZ4压缩,TTL为30天。
## 数据流
1. Producer每秒轮询外部API,将消息收集到队列,然后以批次形式发布到Kafka Topic。
2. 多个Consumer实例分别订阅分配的分区,反序列化JSON,转换时间戳和字段,并将批次写入ClickHouse。
3. 用户通过SQL查询`ecommerce.clickstream_events`表。
## 关键配置
### Producer
- 批次大小:16 KB
- 压缩方式:gzip
- 缓冲内存:32 MB
- 重试次数:3
### Consumer
- 批次大小:1000条事件
- 批次超时:5秒
- 自动提交偏移量已启用
### ClickHouse表
```sql
CREATE TABLE clickstream_events (
user_id String,
session_id String,
timestamp UInt64,
event_type String,
product_id String,
product_category String,
price Float64,
quantity Int32,
source String,
received_time DateTime DEFAULT now()
) ENGINE = MergeTree()
ORDER BY (timestamp, user_id)
```
## 性能报告
- 持续吞吐量:650+ 事件/秒
- 已处理总事件数:510,414+
- 端到端延迟:亚秒级
## 快速入门
1. 使用`uv sync`安装依赖
2. 复制`.env.example`到`.env`并填写Confluent Cloud和ClickHouse凭据
3. 运行Producer:`uv run python services/producer/src/simple_producer.py`
4. 运行Consumer:`uv run python services/consumer/src/simple_consumer.py`
5. 验证数据:`uv run python check_clickhouse.py`
## 故障排除
仓库中包含诊断脚本(`debug_kafka.py`、`test_clickhouse_insert.py`),涵盖常见失败场景,包括API连接超时、Kafka消费者滞后、消息反序列化错误以及ClickHouse插入/模式不匹配问题。
评论
0 评分人数达到10人后显示
登录后参与讨论。