这个项目能做什么

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插入/模式不匹配问题。