这是一个基于Java实现的简易流式计算系统,支持基本的流计算功能和中级特性,包括事件时间窗口处理、数据撤回等功能。
- JDK 11或更高版本
- Maven 3.6或更高版本
src/main/java/com/streamcompute/core/: 核心组件(Graph、Node、Edge等)src/main/java/com/streamcompute/nodes/: 节点实现(Source、Process、Sink节点)src/main/java/com/streamcompute/window/: 窗口处理相关实现src/main/java/com/streamcompute/example/: 示例代码
基础流式系统运行(必须)
- 完整的节点、边处理
- 可持续运行的流式系统
中级流式系统特性(可选)
- 支持事件时间语义下的window处理
- 支持retract处理
- 支持自动故障恢复
高级流式系统特性
- 支持exactly once处理语义
- 克隆项目到本地
- 进入项目目录
- 编译和打包:
mvn clean package
- 运行示例程序:
或者直接运行jar包:
mvn exec:java
java -jar target/simple-stream-compute-1.0-SNAPSHOT-jar-with-dependencies.jar
- 确保已安装JDK 11和Maven
- 打开命令提示符或PowerShell
- 进入项目目录
- 编译和打包:
mvn clean package
- 运行示例程序:
或者直接运行jar包:
mvn exec:java
java -jar target\simple-stream-compute-1.0-SNAPSHOT-jar-with-dependencies.jar
示例程序会创建一个简单的流计算图,包含以下组件:
- 源节点(SourceNode):生成模拟的传感器数据
- 处理节点(ProcessNode):
- 转换数据为可撤回格式
- 进行窗口计算和聚合
- 接收节点(SinkNode):输出计算结果
程序会运行30秒后自动停止,期间会:
- 每500毫秒生成一条传感器数据
- 每1秒触发一次窗口计算
- 随机生成数据撤回操作(10%概率)
程序运行时会输出以下类型的日志:
- 系统启动和停止信息
- 数据处理过程信息
- 窗口触发和计算结果
- 数据撤回操作信息
运行单元测试:
mvn test- 确保系统已安装正确版本的JDK
- 如果遇到内存不足,可以调整JVM参数:
java -Xmx512m -jar target/simple-stream-compute-1.0-SNAPSHOT-jar-with-dependencies.jar
- 程序使用了守护线程,正常情况下会自动退出,如果需要手动停止,可以使用Ctrl+C