Spring AI Kafka事件摘要生成器
一个Spring Boot应用程序,它使用AI从Kafka主题生成实时事件消息的简明摘要。这个演示展示了Spring AI通过模型上下文协议(MCP)与Anthropic的Claude模型和Confluent Kafka的集成。
概述
该应用程序使用来自Kafka主题的消息,其中包含现场事件的观众观察结果,使用AI生成简洁的摘要(最多10个单词),并将这些摘要发布回Kafka以在LED屏幕上显示。
特性
- AI驱动的总结:使用Anthropic的Claude Sonnet 4.5模型生成智能、简洁的摘要
- Kafka集成:通过MCP服务器连接到Confluent Cloud Kafka,用于消费和生成消息
- 聊天记忆:使用消息窗口聊天记忆来维护对话上下文
- 自定义日志记录:包括一个可配置的顾问,用于记录人工智能请求和响应
- 命令行应用程序:作为非web应用程序运行,启动时自动执行
先决条件
- Java 25 或更高
- Gradle (包括包装)
- 无烟煤API密钥:访问Claude AI模型需要
- 融合云访问:
- A. 汇流云 具有活动Kafka集群的帐户 - Kafka集群API证书(在MCP服务器中配置) - 访问 汇流MCP服务器 实例在 https://mcp-confluent.ai-assisted.engineering
安装
- 克隆存储库:
git clone
cd demo- 设置环境变量:
export ANTHROPIC_API_KEY=your_anthropic_api_key_here或者,创建一个 .env 项目根目录中的文件:
ANTHROPIC_API_KEY=your_anthropic_api_key_here配置
应用程序通过以下方式配置 src/main/resources/application.yaml:
AI配置
- 模型:克劳德·十四行诗4.5(
claude-sonnet-4-5-20250929) - API密钥:通过设置
ANTHROPIC_API_KEY环境变量
MCP配置
- MCP服务器URL:
https://mcp-confluent.ai-assisted.engineering - 连接类型:与SSE同步(服务器发送事件)
- 请求超时:60秒
卡夫卡主题
- 输入主题:
user_messages(观众观察) - 输出主题:
llm_summaries(人工智能生成的摘要)
日志记录
- 根级别:警告
- 应用程序级别:DEBUG
- 春季AI级别:信息
- MCP级别:信息
运行应用程序
使用Gradle包装(推荐)
./gradlew bootRun使用Gradle Build
./gradlew build
java -jar build/libs/demo-0.0.1-SNAPSHOT.jar运作原理
- 初创公司:应用程序启动并执行
CommandLineRunner豆 - 消息消费:使用来自的最后5分钟的消息
user_messages卡夫卡主题 - 人工智能处理:向Claude AI发送消息,其中包含生成简明摘要的说明
- 摘要生成:AI创建一个适合家庭的摘要(最多10个字)
- 出版:将摘要发布到
llm_summariesJSON格式的主题,带字段:
- summary:生成的文本 - timestamp:创建摘要时 - messageCount:处理的邮件数
汇流MCP服务器
此应用程序利用 汇流MCP服务器 使AI模型能够通过模型上下文协议(MCP)与Confluent Cloud Kafka集群进行交互。
什么是Confluent MCP服务器?
Confluent MCP Server是一个开源实现,它将AI应用程序与Confluent Cloud连接起来,允许AI模型:
- 消费消息 来自具有可配置时间窗口的Kafka主题
- 生成消息 到Kafka主题
- 列出主题 并检查集群元数据
- 查询模式 从架构注册表
主要特点
- 无缝集成:适用于任何兼容MCP的AI框架(如Spring AI)
- 安全访问:使用Confluent Cloud API密钥进行身份验证
- 实时操作:使AI模型能够实时与流数据交互
- 基于工具的界面:将Kafka操作作为AI模型的可调用工具公开
配置
MCP服务器配置在 application.yaml 要连接到位于的托管实例 https://mcp-confluent.ai-assisted.engineering此服务器实例预先配置了Confluent Cloud凭据,并提供对Kafka集群的访问。
有关更多信息,请访问 .
汇流云
汇流云 是一个完全托管的云原生Apache Kafka服务,提供:
- 完全托管Kafka:不需要基础设施管理
- 全球可用性:跨多个云提供商(AWS、Azure、GCP)部署集群
- 企业安全:内置加密、身份验证和授权
- 模式注册表:数据治理的集中模式管理
- 流处理:集成ksqlDB用于实时数据处理
- 监控和警报:全面的可观察性和警报能力
该应用程序通过MCP服务器连接到Confluent Cloud,在生产级Kafka基础设施上实现AI驱动的操作,而无需管理底层的复杂性。
使用的技术
- 弹簧靴3.5.7:应用程序框架
- 春季AI 1.1.0:人工智能集成框架
- Anthropic Claude:用于文本生成的AI模型
- 模型上下文协议(MCP):AI工具集成协议
- 融合的Kafka:消息流平台
- Gradle:构建自动化工具
- Java 25:编程语言
- JUnit5:测试框架
项目结构
demo/
├── src/
│ ├── main/
│ │ ├── java/com/example/demo/
│ │ │ ├── DemoApplication.java # Main application class
│ │ │ └── SimpleLoggerAdvisor.java # Custom logging advisor
│ │ └── resources/
│ │ ├── application.yaml # Application configuration
│ │ └── logback-spring.xml # Logging configuration
│ └── test/
│ └── java/com/example/demo/
│ └── DemoApplicationTests.java # Unit tests
├── build.gradle.kts # Gradle build configuration
├── settings.gradle.kts # Gradle settings
├── HELP.md # Reference documentation
└── README.md # This file发展
运行测试
./gradlew test建设项目
./gradlew build自定义提示
要修改AI行为,请在中编辑提示 DemoApplication.java:
- 更改摘要长度要求
- 调整音调或风格指南
- 修改内容过滤规则
- 更改消息消费的时间窗口
故障排除
API关键问题
- 确保
ANTHROPIC_API_KEY设置正确 - 验证API密钥是否具有足够的信用和权限
Kafka连接问题
- 检查与MCP服务器的网络连接
- 验证MCP服务器是否可访问
https://mcp-confluent.ai-assisted.engineering - 确保在MCP服务器中配置了正确的Kafka集群凭据
日志记录
- 启用DEBUG日志记录以获取详细信息:设置
logging.level.com.example.demo: DEBUG - 检查
SimpleLoggerAdvisorAI请求/响应详细信息日志
其他资源
许可证
此项目根据MIT许可证获得许可-请参阅 许可证 文件以获取详细信息。
版本
版本:0.0.1快照
