健康堆栈数据平台
用于实时健康数据收集、处理和分析的端到端数据平台。通过Kafka收集iOS应用程序收集的健康数据,通过Apache Flink实时处理并存储在Apache Iceberg数据湖中,AI代理可以通过MCP服务器查询数据。
系统体系结构
┌─────────────┐
│ iOS App │
└──────┬──────┘
│ JSON/HTTPS
▼
┌──────────────────────────────────────────┐
│ API Gateway (FastAPI) │
│ ┌──────────┐ ┌────────────────┐ │
│ │Validator │→ │Avro Converter │ │
│ └──────────┘ └────────┬───────┘ │
└─────────────────────────┼────────────────┘
│
▼
┌───────────────────────┐
│ Schema Registry │
└───────────────────────┘
│
▼
┌───────────────────────┐
│ Kafka Cluster │
│ (3 brokers) │
│ health-data-raw │
└───────────┬───────────┘
│
▼
┌──────────────────────────────────────────┐
│ Apache Flink Cluster │
│ ┌────────────────────────────────────┐ │
│ │ Stream Processing Pipeline │ │
│ │ - Transformation │ │
│ │ - Validation │ │
│ │ - Time-based Aggregation │ │
│ │ (Daily/Weekly/Monthly) │ │
│ └────────────┬───────────────────────┘ │
└───────────────┼──────────────────────────┘
│
▼
┌──────────────────────────────────────────┐
│ Apache Iceberg Data Lake │
│ ┌──────────────────────────────────┐ │
│ │ health_data_raw │ │
│ │ health_data_daily_agg │ │
│ │ health_data_weekly_agg │ │
│ │ health_data_monthly_agg │ │
│ └──────────────────────────────────┘ │
│ │
│ Storage: MinIO (S3-compatible) │
└──────────────┬───────────────────────────┘
│
▼
┌──────────────────────────────────────────┐
│ Health Data MCP Server │
│ ┌──────────────────────────────────┐ │
│ │ Query Tools for AI Agents │ │
│ │ - get_daily_aggregates │ │
│ │ - get_weekly_aggregates │ │
│ │ - get_monthly_aggregates │ │
│ │ - get_top_records │ │
│ └──────────────────────────────────┘ │
└──────────────┬───────────────────────────┘
│
▼
┌──────────┐
│ AI Agent │
│ (Kiro) │
└──────────┘主要元件
1.API网关(FastAPI)
接收iOS应用程序发送的健康数据并将其传递到Kafka的REST API网关。
主要功能:
- 验证和转换JSON数据
- 以Avro格式序列化
- Schema Registry集成
- 发布Kafka消息
- 错误处理和DLQ(Dead Letter Queue)
技术堆栈:
- FastAPI(Python 3.11+)
- Kafka与Python的融合
- Avro架构注册表
- Pydantic数据验证
端口:
- APIhttp://localhost:3000
- Swagger用户界面:http://localhost:3000/docs
- 健康检查:http://localhost:3000/health
📖 详细文档: gateway/README.md
______________________________________________________________________
2.Kafka集群
通过高可用性消息代理可靠地传输健康数据流。
配置:
- 3个Kafka代理(KRaft模式)
- Schema Registry(Avro模式管理)
- Kafka UI(监视和管理)
主题:
health-data-raw:原始健康数据(6分区,RF=3)health-data-dlq:失败的消息(3分区,RF=3)
端口:
- 代理1:本地主机:19092
- 经纪人2:本地主机:19093
- 经纪人3:本地主机:19094
- 架构注册表:http://localhost:8081
- Kafka用户界面:http://localhost:8080
______________________________________________________________________
3.Flink消费者(Apache Flink)
作为实时流处理应用程序,Kafka将消耗数据并将其存储在Iceberg中。
主要功能:
- 实时数据转换和验证
- 基于时间的统计(每日/每周/每月)
- Exactly-once保证浪漫
- 晚点数据处理(Late Data Handling)
- 基于检查点的故障转移
统计:
- min_value、max_value、avg_value
- sum_value、count、stddev_value
- first_value,last_value
技术堆栈:
- Apache Flink 1.18+
- PyFlink(Python API)
- 阿帕奇冰山
- MinIO(S3兼容存储)
端口:
- Flink Web用户界面:http://localhost:8081
- 普罗米修斯指标:http://localhost:9249
📖 详细文档: flink_consumer/README.md
______________________________________________________________________
4.阿帕奇冰山数据湖
利用可扩展的数据湖高效存储和查询健康数据。
表格:
health_data_raw:原始健康数据health_data_daily_agg:每日统计health_data_weekly_agg:每周统计health_data_monthly_agg:月度统计health_data_errors:错误日志(DLQ)
分区:
user_id(哈希分区)aggregation_date(日期分区)data_type(类别分区)
存储:
- MinIO(兼容S3)
- 仓库:s3a://数据湖/仓库
- 检查点:s3a://flink检查点
端口:
- API矿业公司:http://localhost:9000
- MinIO控制台:http://localhost:9001(最小负载/最小负载)
______________________________________________________________________
5.健康数据MCP服务器
AI代理是一个模型上下文协议(MCP)服务器,可用于查看Iceberg数据湖的健康数据。
提供工具:
get_daily_aggregates:日间统计查询get_weekly_aggregates:查询每周统计get_monthly_aggregates:月度统计查询get_top_records:查看最高/最低记录
技术堆栈:
- MCP SDK(Python)
- 皮冰山
- PyArrow
使用示例:
# Kiro AI 에이전트에서 사용
"최근 30일간 user-123의 심박수 평균을 알려줘"
→ get_daily_aggregates(user_id="user-123", data_type="heartRate")
"이번 달 가장 많이 걸었던 날은?"
→ get_top_records(user_id="user-123", data_type="steps", sort_by="sum_value")📖 详细文档: health_data_mcp/README.md
______________________________________________________________________
快速入门
前提条件
- Docker&Docker编写
- Python3.11+(本地开发时)
- uv(Python包管理器)
# uv 설치
curl -LsSf https://astral.sh/uv/install.sh | sh运行整个堆栈
# 1. 모든 서비스 시작
docker-compose up -d
# 2. 서비스 상태 확인
docker-compose ps
# 3. 로그 확인
docker-compose logs -f
# 4. 서비스 접속
# - API Gateway: http://localhost:3000/docs
# - Kafka UI: http://localhost:8080
# - Flink UI: http://localhost:8081
# - MinIO Console: http://localhost:9001测试数据传输
# 샘플 헬스 데이터 전송
curl -X POST http://localhost:3000/api/v1/health-data \
-H "Content-Type: application/json" \
-d '{
"userId": "user-123",
"timestamp": "2025-11-26T10:30:00Z",
"dataType": "heart_rate",
"value": 72,
"unit": "bpm",
"metadata": {
"deviceId": "iPhone14-ABC123",
"appVersion": "1.2.3",
"platform": "iOS"
}
}'验证数据
# 1. Kafka UI에서 메시지 확인
# http://localhost:8080 → Topics → health-data-raw
# 2. Flink UI에서 처리 상태 확인
# http://localhost:8081 → Jobs
# 3. MinIO에서 Iceberg 파일 확인
# http://localhost:9001 → data-lake → warehouse______________________________________________________________________
MCP服务器设置(Kiro)
1.安装MCP服务器
cd health_data_mcp
uv sync2.添加Kiro配置文件
.kiro/settings/mcp.json 将以下内容添加到文件:
{
"mcpServers": {
"health-data": {
"command": "python",
"args": ["-m", "health_data_mcp.main"],
"cwd": "/path/to/health_data_mcp",
"env": {
"ICEBERG_CATALOG_URI": "http://localhost:8181",
"ICEBERG_CATALOG_NAME": "health_catalog",
"ICEBERG_WAREHOUSE": "s3://data-lake/warehouse",
"ICEBERG_DATABASE": "health_data",
"S3_ENDPOINT": "http://localhost:9000",
"S3_ACCESS_KEY": "minioadmin",
"S3_SECRET_KEY": "minioadmin"
},
"disabled": false
}
}
}3.在Kiro中使用
"user-123의 최근 30일 심박수 데이터를 분석해줘"
"이번 주 가장 많이 걸었던 날은?"
"지난 6개월 월별 평균 심박수 추이를 보여줘"______________________________________________________________________
数据流
1.数据收集(Ingestion)
iOS App → API Gateway → Kafka (health-data-raw)2.实时处理(Processing)
Kafka → Flink Consumer → Transformation → Validation3.聚合(Aggregation)
Flink → Time Windows (Daily/Weekly/Monthly) → Statistics4.存储(Storage)
Flink → Iceberg Tables → MinIO (S3)5.查询(Query)
AI Agent → MCP Server → PyIceberg → Iceberg Tables______________________________________________________________________
支持的数据类型
数据类型说明单位示例值 |------------|------|------|---------| | heart_rate 心率bpm72 | steps 步数count8543 | distance 移动距离公里5.2 | blood_pressure | 혈압 | mmHg |{收缩压:120,舒张压:80}| | blood_glucose 血糖mg/dL 95 | body_temperature 体温°C 36.8 | oxygen_saturation 氧饱和度%98| | respiratory_rate 呼吸数breaths/min 16 | weight 体重70.5公斤 | sleep |睡眠| minutes|{duration:480,…}|
______________________________________________________________________
监控和管理
检查服务状态
# 모든 서비스 상태
docker-compose ps
# 특정 서비스 로그
docker-compose logs -f gateway
docker-compose logs -f flink-jobmanager
docker-compose logs -f kafka-broker-1Web UI连接
服务URL说明 |--------|-----|------| |API网关|http://localhost:3000/docs| Swagger用户界面| | Kafka UI | http://localhost:8080| Kafka主题和信息| | Flink UI | http://localhost:8081|监控Flink任务| | MinIO Console | http://localhost:9001| S3存储管理|
健康检查
# API Gateway
curl http://localhost:3000/health
# Flink JobManager
curl http://localhost:8081/overview
# Schema Registry
curl http://localhost:8081/subjects收集度量
# API Gateway 메트릭 (Prometheus)
curl http://localhost:3000/metrics
# Flink 메트릭
curl http://localhost:9249/metrics______________________________________________________________________
设置开发环境
Gateway本地开发
cd gateway
uv sync
cp .env.example .env
# .env 파일 수정
uv run uvicorn main:app --reload --host 0.0.0.0 --port 3000Flink Consumer本地开发
cd flink_consumer
uv sync
cp .env.example .env.local
# .env.local 파일 수정
source .venv/bin/activate
python main.pyMCP Server本地开发
cd health_data_mcp
uv sync
cp .env.example .env
# .env 파일 수정
python -m health_data_mcp.main______________________________________________________________________
测试
Gateway测试
cd gateway
uv run pytestFlink Consumer测试
cd flink_consumer
uv run pytest测试MCP Server
cd health_data_mcp
uv run pytest集成测试
# 전체 스택 통합 테스트
./run_integration_tests.sh______________________________________________________________________
性能和可扩展性
吞吐量(Throughput)
- API网关:~ 10000 req/s(单实例)
- 卡夫卡:~ 100000 msg/秒(3中间人)
- 弗林克:~ 50000 records/sec(12个任务插槽)
扩展方法
网关水平扩展:
docker-compose up -d --scale gateway=3Flink TaskManager扩展:
docker-compose up -d --scale flink-taskmanager=5Kafka分区增加:
docker-compose exec kafka-broker-1 kafka-topics \
--alter --topic health-data-raw \
--partitions 12 \
--bootstrap-server kafka-broker-1:9092______________________________________________________________________
故障转移
Flink检查点
- 每隔60秒自动检查点
- 将状态保存到S3(MinIO)
- 发生故障时自动恢复
Kafka复制
- 将数据复制到3个代理(RF=3)
- 至少同步两个代理(min.insync.replicas=2)
- 代理出现故障时自动进行故障切换
数据完整性
- Exactly-once保证浪漫
- 基于事务的Kafka制作人
- 基于Flink检查点的状态管理
______________________________________________________________________
故障排除
服务未启动时
# 로그 확인
docker-compose logs
# 서비스 재시작
docker-compose restart
# 전체 재시작
docker-compose down
docker-compose up -dKafka连接错误
# Kafka 브로커 상태 확인
docker-compose exec kafka-broker-1 kafka-broker-api-versions \
--bootstrap-server kafka-broker-1:9092
# 토픽 목록 확인
docker-compose exec kafka-broker-1 kafka-topics \
--list --bootstrap-server kafka-broker-1:9092Flink操作失败
# Flink 로그 확인
docker-compose logs flink-jobmanager
docker-compose logs flink-taskmanager-1
# Flink UI에서 작업 상태 확인
# http://localhost:8081详细的故障排除指南
______________________________________________________________________
生产部署
设置环境变量
# 프로덕션 환경 변수 설정
LOG_LEVEL=INFO
ENVIRONMENT=production安全设置
# API Gateway HTTPS 활성화
SSL_ENABLED=true
SSL_CERTFILE=/path/to/cert.pem
SSL_KEYFILE=/path/to/key.pem
# Kafka SASL 인증 (선택사항)
KAFKA_SECURITY_PROTOCOL=SASL_SSL
KAFKA_SASL_MECHANISM=SCRAM-SHA-512Kubenetes部署
# Flink Operator 설치
kubectl apply -f flink_consumer/k8s/
# Gateway 배포
kubectl apply -f gateway/k8s/______________________________________________________________________
文件
特定于组件的详细文档
API文档
操作指南
体系结构文档
______________________________________________________________________
技术堆栈
后端
- Python 3.11+ -主要编程语言
- 快速API -API Gateway框架
- Apache Flink 1.18+ -流处理
- Apache Kafka 7.5 -消息代理
- 阿帕奇冰山 -数据湖表格式
存储
- MinIO -S3-compatible对象存储
- 冲突模式注册表 -Avro模式管理
监控
- 普罗米修斯 -收集度量
- 格拉法纳 -公制可视化(可选)
- Kafka用户界面 -Kafka监控
开发工具
- 紫外线 -Python包管理器
- Docker&Docker编写 -集装箱化
- pytest -测试框架
- 颈毛 -林特和波马特
______________________________________________________________________
许可证
专有-健康堆栈项目
______________________________________________________________________
贡献
欢迎焦点和全员任务!
______________________________________________________________________
支持
如果遇到问题,请检查以下内容:
- 故障排除指南
- 每个元件的README文件
- 日志文件(
docker-compose logs)
______________________________________________________________________
上次更新: 2025-11-26
