# doit36_realtime_project **Repository Path**: edidada/doit36_realtime_project ## Basic Information - **Project Name**: doit36_realtime_project - **Description**: 多易教育实时数仓36期项目课代码 - **Primary Language**: Unknown - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 5 - **Created**: 2026-03-18 - **Last Updated**: 2026-03-19 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # doit36_realtime_project ## 1. 项目介绍 本项目是多易大数据的Java实时项目,主要包含以下模块: - **realtime_dw**: 实时数仓模块,基于Flink进行实时数据处理和ETL - **batch_jobs**: 批处理任务模块,基于Spark进行离线数据处理 - **redis_demo**: Redis操作示例模块 - **rt_rule_mk**: 实时规则引擎模块,基于Flink实现动态规则匹配 项目主要功能包括: - 实时用户行为日志处理和维度打宽 - 实时流量概况多维模型统计 - 页面访问时长分析 - 搜索事件分析 - 视频事件聚合 - 广告请求实时特征工程 - 实时规则匹配和用户画像 - 数据同步(MySQL CDC、HBase等) ## 2. 项目环境 ### 开发环境要求 - **Java**: 17+ (当前环境: OpenJDK 11.0.21,项目配置为Java 8) - **Maven**: 3.6.1+ - **Scala**: 2.12.12 (用于batch_jobs模块) ### 运行时依赖环境 - **Kafka**: 消息队列,用于实时数据传输 - **MySQL**: 数据库,用于存储业务数据和规则元数据 - **HBase**: NoSQL数据库,用于存储用户画像和标签 - **Redis**: 缓存数据库 - **Doris**: OLAP数据库,用于实时数据分析 - **Hadoop**: 分布式存储和计算框架 ## 3. 项目依赖 ### 主要依赖版本 #### Flink相关 (realtime_dw, rt_rule_mk) - Flink: 1.16.1 - Flink Connector Kafka: 1.16.1 - Flink Connector HBase: 1.16.1 - Flink Connector JDBC: 1.16.1 - Flink StateBackend RocksDB: 1.16.1 - Flink CEP: 1.16.1 - Flink Table Planner: 1.16.1 #### Spark相关 (batch_jobs) - Spark: 3.1.2 - Spark SQL: 3.1.2 - Spark Hive: 3.1.2 - Scala: 2.12.12 #### 数据库连接 - MySQL Connector: 8.0.30 - HBase Client: 2.4.9 - Redis Jedis: 4.3.1 - Elasticsearch: 7.17.5 #### 工具库 - Lombok: 1.18.24 - FastJSON: 1.2.83 - Geohash: 1.4.0 - RoaringBitmap: 0.9.39 - Groovy: 3.0.12 - HttpClient: 4.5.14 #### CDC相关 - Flink Connector MySQL CDC: 2.3.0 - Flink Doris Connector: 1.1.0 #### 日志 - SLF4J Log4j12: 1.7.30 - Log4j: 1.2.16 #### 测试 - JUnit: 4.13.1 ## 4. 项目配置 ### Maven配置 项目采用多模块Maven结构,父POM定义了公共配置: ```xml cn.doitedu doit36_realtime_project 1.0-SNAPSHOT pom ``` 子模块: - realtime_dw - batch_jobs - redis_demo - rt_rule_mk ### 日志配置 使用Log4j进行日志管理,配置文件位于: - `realtime_dw/src/main/resources/log4j.properties` - `rt_rule_mk/src/main/resources/log4j.properties` 日志级别:INFO 输出格式:`[%-5p] %d --> %m %x %n` ### Flink配置 - 并行度:默认为1 - Checkpoint间隔:2000-5000ms - Checkpoint模式:EXACTLY_ONCE - State Backend:HashMapStateBackend / EmbeddedRocksDBStateBackend - Checkpoint存储:`file:/d:/ckpt` 或 `file:/d:/checkpoint` ### Kafka配置 - Bootstrap Servers: `doitedu:9092` - Topics: `dwd-events-log` (行为日志明细宽表) - 消费组:动态生成 ### 数据库配置 - MySQL: `doitedu:3306`, 用户名: `root`, 密码: `root`, 数据库: `realtimedw` - HBase: `doitedu:2181` (ZooKeeper) - Redis: `doitedu:6379` ## 5. 项目运行 ### 编译项目 ```bash # 清理并编译所有模块 mvn clean compile # 打包项目 mvn clean package # 跳过测试打包 mvn clean package -DskipTests ``` ### 运行Flink作业 #### realtime_dw模块 ```bash # 运行用户行为日志维度打宽作业 java -cp realtime_dw/target/realtime_dw-1.0-SNAPSHOT.jar cn.doitedu.rtdw.data_etl.EtlJob01_UserEventsLogCommonDim # 运行实时流量概况多维模型统计 java -cp realtime_dw/target/realtime_dw-1.0-SNAPSHOT.jar cn.doitedu.rtdw.data_etl.EtlJob02_MallTfcAg01 # 运行页面访问时长分析 java -cp realtime_dw/target/realtime_dw-1.0-SNAPSHOT.jar cn.doitedu.rtdw.data_etl.EtlJob03_PageAccessTimeLong # 运行搜索事件分析 java -cp realtime_dw/target/realtime_dw-1.0-SNAPSHOT.jar cn.doitedu.rtdw.data_etl.EtlJob06_SearchEventsAnalyse # 运行视频事件聚合 java -cp realtime_dw/target/realtime_dw-1.0-SNAPSHOT.jar cn.doitedu.rtdw.data_etl.EtlJob05_VideoEventsAgg # 运行广告请求实时特征工程 java -cp realtime_dw/target/realtime_dw-1.0-SNAPSHOT.jar cn.doitedu.rtdw.data_etl.EtlJob07_AdRequestSessionProcess ``` #### rt_rule_mk模块 ```bash # 运行实时规则引擎Demo java -cp rt_rule_mk/target/rt_rule_mk-1.0-SNAPSHOT.jar cn.doitedu.rtmk.demo1.Demo1 java -cp rt_rule_mk/target/rt_rule_mk-1.0-SNAPSHOT.jar cn.doitedu.rtmk.demo2.MultiRuleEngine java -cp rt_rule_mk/target/rt_rule_mk-1.0-SNAPSHOT.jar cn.doitedu.rtmk.demo3.Demo3 java -cp rt_rule_mk/target/rt_rule_mk-1.0-SNAPSHOT.jar cn.doitedu.rtmk.demo4.Demo4 java -cp rt_rule_mk/target/rt_rule_mk-1.0-SNAPSHOT.jar cn.doitedu.rtmk.demo5.Demo5 java -cp rt_rule_mk/target/rt_rule_mk-1.0-SNAPSHOT.jar cn.doitedu.rtmk.demo6.Demo6 java -cp rt_rule_mk/target/rt_rule_mk-1.0-SNAPSHOT.jar cn.doitedu.rtmk.demo7.Demo7 java -cp rt_rule_mk/target/rt_rule_mk-1.0-SNAPSHOT.jar cn.doitedu.rtmk.demo8.Demo8 java -cp rt_rule_mk/target/rt_rule_mk-1.0-SNAPSHOT.jar cn.doitedu.rtmk.demo9.Demo9 ``` #### batch_jobs模块 ```bash # 运行批处理作业 java -cp batch_jobs/target/batch_jobs-1.0-SNAPSHOT.jar cn.doitedu.batch.jobs.DataLoadJob01_GeoAreaInfo2Hbase ``` #### redis_demo模块 ```bash # 运行Redis示例 java -cp redis_demo/target/redis_demo-1.0-SNAPSHOT.jar cn.doitedu.jedis.RedisDemo ``` ### Flink集群运行 ```bash # 提交作业到Flink集群 flink run -c cn.doitedu.rtdw.data_etl.EtlJob01_UserEventsLogCommonDim realtime_dw/target/realtime_dw-1.0-SNAPSHOT.jar # 提交作业并指定并行度 flink run -p 4 -c cn.doitedu.rtdw.data_etl.EtlJob01_UserEventsLogCommonDim realtime_dw/target/realtime_dw-1.0-SNAPSHOT.jar ``` ## 6. 项目测试 ### 单元测试 ```bash # 运行所有测试 mvn test # 运行特定模块的测试 mvn test -pl rt_rule_mk # 运行特定测试类 mvn test -Dtest=TestClassName ``` ### 测试数据 项目包含多个测试数据文件: - `model1_test.json`: 规则模型1测试数据 - `model2_test.json`: 规则模型2测试数据 - `model3_test.json`: 规则模型3测试数据 - `model4_test.json`: 规则模型4测试数据 - `model5_test.json`: 规则模型5测试数据 - `model5_test2.json`: 规则模型5测试数据2 - `x.json`: 其他测试数据 ### 测试类 - `cn.doitedu.test.EventsMoni`: 事件监控测试 - `cn.doitedu.test.HbaseProfileTagsMoni`: HBase用户画像标签测试 - `cn.doitedu.test.HistoryConditionValueMoni`: 历史条件值测试 - `cn.doitedu.test.MysqlRuleMetaMoni`: MySQL规则元数据测试 ### 监控和调试 - Flink Web UI: `http://localhost:8081` (默认端口) - 日志输出: 控制台输出,格式为Log4j配置的格式 - Checkpoint监控: 可通过Flink Web UI查看Checkpoint状态 ### 注意事项 1. 确保所有依赖服务(Kafka、MySQL、HBase、Redis等)已启动 2. 检查配置文件中的连接地址和端口是否正确 3. 首次运行前需要创建必要的数据库表和Kafka topics 4. 建议先在本地模式测试,确认无误后再提交到集群 5. 注意Checkpoint存储路径的配置,确保有足够的磁盘空间 6. 生产环境建议使用RocksDB State Backend以获得更好的性能 ## 项目结构 ``` doit36_realtime_project/ ├── realtime_dw/ # 实时数仓模块 │ ├── src/main/java/cn/doitedu/rtdw/ │ │ ├── beans/ # 数据Bean定义 │ │ ├── data_etl/ # ETL作业 │ │ ├── data_sync/ # 数据同步作业 │ │ ├── rt_dashboard/ # 实时看板 │ │ └── utils/ # 工具类 │ └── pom.xml ├── batch_jobs/ # 批处理模块 │ ├── src/main/java/cn/doitedu/batch/jobs/ │ └── pom.xml ├── redis_demo/ # Redis示例模块 │ ├── src/main/java/cn/doitedu/jedis/ │ └── pom.xml ├── rt_rule_mk/ # 实时规则引擎模块 │ ├── src/main/java/cn/doitedu/rtmk/ │ │ ├── beans/ # 规则Bean │ │ ├── demo1-9/ # 各版本演示 │ │ └── groovy/ # Groovy规则脚本 │ └── pom.xml ├── lib/ # 第三方jar包 │ └── flink-shaded-hadoop-3-uber-3.1.1.7.2.9.0-173-9.0.jar ├── pom.xml # 父POM └── README.md # 项目说明文档 ``` ## 技术栈总结 - **实时计算**: Apache Flink 1.16.1 - **批处理**: Apache Spark 3.1.2 - **消息队列**: Apache Kafka - **数据库**: MySQL, HBase, Redis, Doris - **搜索引擎**: Elasticsearch 7.17.5 - **编程语言**: Java 8+, Scala 2.12.12 - **构建工具**: Maven 3.6.1 - **日志**: Log4j + SLF4J - **其他**: Groovy 3.0.12, RoaringBitmap, Geohash