Flink实战:5分钟搞定城市交通卡口超速监控(附完整代码)
Flink实战:5分钟搞定城市交通卡口超速监控(附完整代码)
最近在和一些做智慧城市项目的朋友聊天,他们普遍提到一个痛点:交通流数据上来了,但实时分析能力跟不上。摄像头每秒都在产生海量的过车记录,如何从中快速、准确地揪出那些超速的“马路飞车”,成了摆在技术团队面前的一道坎。传统的批处理方案延迟太高,等报表出来,车早就跑没影了;自己写流处理程序,又面临着状态管理、数据一致性、系统容错等一系列复杂问题。这时候,一个成熟、高效的流处理框架就显得至关重要。
Apache Flink,作为当前流计算领域的“当红炸子鸡”,其核心优势就在于真正的流式处理和精确一次(Exactly-Once)的状态一致性保证。这意味着,用它来处理像交通卡口监控这样的实时场景,不仅速度快,而且结果可靠,不会因为系统故障或网络波动而漏掉或重复计算违规车辆。今天,我们就抛开复杂的理论,直接上手,用Flink构建一个极简但五脏俱全的实时超速监控原型。我会把每一步的代码都贴出来,并解释关键设计思路,目标是让你在理解原理的同时,能快速复现出一个可运行的系统。
1. 环境准备与项目初始化
在开始敲代码之前,我们需要把“厨房”收拾好。这里假设你已经有了Java和Maven的基本环境。我们的技术栈选择是:Flink 1.14+(API相对稳定),用Scala语言编写(兼顾表达力和性能),数据源模拟Kafka,结果存到MySQL。别被这个列表吓到,我们一步步来。
首先,打开你的IDE(IntelliJ IDEA或VS Code都可以),创建一个Maven项目。项目的pom.xml文件是依赖管理的核心,我们需要在这里声明所有必要的库。
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.example</groupId>
<artifactId>flink-traffic-monitor</artifactId>
<version>1.0-SNAPSHOT</version>
<properties>
<flink.version>1.14.4</flink.version>
<scala.binary.version>2.12</scala.binary.version>
<kafka.version>2.8.0</kafka.version>
<mysql.version>8.0.28</mysql.version>
<maven.compiler.source>8</maven.compiler.source>
<maven.compiler.target>8</maven.compiler.target>
</properties>
<dependencies>
<!-- Flink核心依赖 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-scala_${scala.binary.version}</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-clients_${scala.binary.version}</artifactId>
<version>${flink.version}</version>
</dependency>
<!-- Kafka连接器,用于消费数据 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-kafka_${scala.binary.version}</artifactId>
<version>${flink.version}</version>
</dependency>
<!-- JDBC连接器,用于写入MySQL -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-jdbc_${scala.binary.version}</artifactId>
<version>${flink.version}</version>
</dependency>
<!-- MySQL驱动 -->
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
<version>${mysql.version}</version>
</dependency>
<!-- 日志 -->
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-simple</artifactId>
<version>1.7.36</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>net.alchim31.maven</groupId>
<artifactId>scala-maven-plugin</artifactId>
<version>4.7.1</version>
<executions>
<execution>
<goals>
<goal>compile</goal>
<goal>testCompile</goal>
</goals>
</execution>
</executions>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.2.4</version>
<executions>
<execution>
<phase>package</phase>
<goals>
<goal>shade</goal>
</goals>
<configuration>
<filters>
<filter>
<artifact>*:*</artifact>
<excludes>
<exclude>META-INF/*.SF</exclude>
<exclude>META-INF/*.DSA</exclude>
<exclude>META-INF/*.RSA</exclude>
</excludes>
</filter>
</filters>
</configuration>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>
注意:Flink版本建议选择1.14或更高,这个版本的API已经非常成熟稳定。Scala版本与Flink版本有对应关系,这里选用2.12。
依赖搞定后,我们在src/main/scala目录下创建我们的第一个Scala对象:TrafficSpeedMonitor。同时,我们需要定义几个核心的数据结构(Case Class),它们就像是流中数据的“形状”。
// 定义数据结构的对象
object TrafficModels {
// 原始卡口过车数据
case class TrafficLog(
actionTime: Long, // 时间戳,毫秒
monitorId: String, // 卡口ID,如 "0001"
cameraId: String, // 摄像头ID
carPlate: String, // 车牌号,如 "京A12345"
speed: Double, // 瞬时速度,单位 km/h
roadId: String, // 道路ID
areaId: String // 区域ID
)
// 卡口限速信息(来自维度表)
case class SpeedLimit(
monitorId: String,
roadId: String,
limitSpeed: Int // 限速值,单位 km/h
)
// 超速告警结果
case class SpeedingViolation(
carPlate: String,
monitorId: String,
roadId: String,
actualSpeed: Double,
limitSpeed: Int,
actionTime: Long,
processTime: Long = System.currentTimeMillis() // 处理时间
)
}
这些case class不仅仅是数据的容器,它们在Flink的类型系统中也扮演着重要角色,能帮助Flink更高效地进行序列化和反序列化。接下来,我们进入核心逻辑部分。
2. 构建实时数据流与广播状态模式
超速判断的核心逻辑其实不复杂:将实时流过的车辆速度,与对应卡口的限速标准进行比对。难点在于,限速标准这个“规则”可能发生变化(比如某路段施工临时降速),而且它相对于海量的车辆数据流,属于“小数据”。Flink的广播状态(Broadcast State) 模式正是为这种“大流小表”的关联场景量身定做的。
它的工作原理可以想象成一个广播电台和一个听众群体。限速规则表是“电台”,它把更新的规则(比如“卡口0001限速调整为60”)广播给所有“听众”(即处理车辆数据的算子实例)。每个听众都有一份最新的规则手册,当车辆数据流过时,就能立刻用最新规则进行判断。
首先,我们需要在MySQL中创建对应的表:
-- 卡口限速维度表
CREATE TABLE `speed_limit_info` (
`id` int NOT NULL AUTO_INCREMENT,
`monitor_id` varchar(20) NOT NULL,
`road_id` varchar(20) NOT NULL,
`limit_speed` int NOT NULL COMMENT '限速值(km/h)',
`update_time` timestamp NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
PRIMARY KEY (`id`),
UNIQUE KEY `idx_monitor_road` (`monitor_id`,`road_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
-- 插入一些示例数据
INSERT INTO `speed_limit_info` (monitor_id, road_id, limit_speed) VALUES
('0001', 'R001', 60),
('0002', 'R001', 80),
('0003', 'R002', 100),
('0004', 'R003', 40);
-- 超速告警结果表
CREATE TABLE `speeding_violation` (
`id` bigint NOT NULL AUTO_INCREMENT,
`car_plate` varchar(20) NOT NULL,
`monitor_id` varchar(20) DEFAULT NULL,
`road_id` varchar(20) DEFAULT NULL,
`actual_speed` decimal(5,1) DEFAULT NULL,
`limit_speed` int DEFAULT NULL,
`action_time` bigint DEFAULT NULL COMMENT '过车时间戳',
`process_time` bigint DEFAULT NULL COMMENT '系统处理时间戳',
`created_at` timestamp NULL DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (`id`),
KEY `idx_car_time` (`car_plate`,`action_time`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
现在,让我们编写Flink作业的主类。代码虽长,但结构清晰,我加了详细注释。
import org.apache.flink.api.common.state.{BroadcastState, MapStateDescriptor, ReadOnlyBroadcastState}
import org.apache.flink.configuration.Configuration
import org.apache.flink.streaming.api.functions.co.BroadcastProcessFunction
import org.apache.flink.streaming.api.functions.source.{RichSourceFunction, SourceFunction}
import org.apache.flink.streaming.api.scala._
import org.apache.flink.util.Collector
import java.sql.{Connection, DriverManager, PreparedStatement, ResultSet}
import java.util.concurrent.TimeUnit
object TrafficSpeedMonitor {
import TrafficModels._
// 定义广播状态的描述符,用于在算子间传递限速规则
private final val SPEED_LIMIT_DESCRIPTOR = new MapStateDescriptor[String, Int](
"speedLimitState",
classOf[String], // Key的类型:用 "monitorId:roadId" 作为唯一键
classOf[Int] // Value的类型:限速值
)
def main(args: Array[String]): Unit = {
// 1. 创建流执行环境
val env = StreamExecutionEnvironment.getExecutionEnvironment
env.setParallelism(1) // 本地测试设为1,生产环境根据资源调整
// 2. 模拟车辆数据源(实际生产环境应连接Kafka)
val trafficLogStream: DataStream[TrafficLog] = env
.addSource(new SimulatedTrafficSource)
.name("simulated-traffic-source")
// 3. 广播流:定期从MySQL读取限速规则
val speedLimitBroadcastStream: BroadcastStream[SpeedLimit] = env
.addSource(new JdbcSpeedLimitSource(5)) // 每5秒查询一次
.broadcast(SPEED_LIMIT_DESCRIPTOR)
// 4. 连接主流与广播流,进行超速判断
val speedingViolationStream: DataStream[SpeedingViolation] = trafficLogStream
.connect(speedLimitBroadcastStream)
.process(new SpeedCheckBroadcastProcessFunction)
.name("speed-check-processor")
// 5. 将超速告警写入MySQL
speedingViolationStream.addSink(new JdbcSinkSpeedingViolation)
.name("jdbc-sink-violation")
// 6. 同时,我们也可以打印到控制台方便调试
speedingViolationStream.print().name("console-sink")
// 7. 执行作业
env.execute("Real-time Traffic Speed Monitoring")
}
// --- 以下是内部类实现 ---
/**
* 模拟交通卡口数据源。
* 在实际项目中,这里应该替换为从Kafka消费真实的过车数据。
*/
class SimulatedTrafficSource extends RichSourceFunction[TrafficLog] {
@volatile private var isRunning = true
private val rand = new scala.util.Random
override def run(ctx: SourceFunction.SourceContext[TrafficLog]): Unit = {
val monitorIds = Array("0001", "0002", "0003", "0004")
val roadIds = Array("R001", "R002", "R003")
val areaIds = Array("01", "02", "03")
val carPrefixes = Array("京A", "京B", "沪A", "粤S")
while (isRunning) {
// 模拟生成一条过车记录
val log = TrafficLog(
actionTime = System.currentTimeMillis() - rand.nextInt(5000), // 模拟5秒内的时间
monitorId = monitorIds(rand.nextInt(monitorIds.length)),
cameraId = f"CAM${rand.nextInt(10000)}%05d",
carPlate = carPrefixes(rand.nextInt(carPrefixes.length)) + f"${rand.nextInt(10000)}%05d",
speed = 30 + rand.nextDouble() * 100, // 速度在30-130 km/h之间
roadId = roadIds(rand.nextInt(roadIds.length)),
areaId = areaIds(rand.nextInt(areaIds.length))
)
ctx.collect(log)
TimeUnit.MILLISECONDS.sleep(100) // 每100毫秒模拟一辆车
}
}
override def cancel(): Unit = {
isRunning = false
}
}
/**
* 从MySQL定期拉取限速规则的Source。
* 这是一个RichSourceFunction,可以方便地管理数据库连接。
*/
class JdbcSpeedLimitSource(intervalSec: Int) extends RichSourceFunction[SpeedLimit] {
@volatile private var isRunning = true
private var connection: Connection = _
private var statement: PreparedStatement = _
override def open(parameters: Configuration): Unit = {
Class.forName("com.mysql.cj.jdbc.Driver")
connection = DriverManager.getConnection(
"jdbc:mysql://localhost:3306/traffic_db?useSSL=false&serverTimezone=UTC",
"your_username",
"your_password"
)
statement = connection.prepareStatement("SELECT monitor_id, road_id, limit_speed FROM speed_limit_info")
}
override def run(ctx: SourceFunction.SourceContext[SpeedLimit]): Unit = {
while (isRunning) {
val rs: ResultSet = statement.executeQuery()
while (rs.next()) {
val limit = SpeedLimit(
monitorId = rs.getString("monitor_id"),
roadId = rs.getString("road_id"),
limitSpeed = rs.getInt("limit_speed")
)
ctx.collect(limit)
}
rs.close()
TimeUnit.SECONDS.sleep(intervalSec) // 定期更新
}
}
override def cancel(): Unit = {
isRunning = false
if (statement != null) statement.close()
if (connection != null) connection.close()
}
override def close(): Unit = {
cancel()
}
}
/**
* 核心处理逻辑:BroadcastProcessFunction。
* 它有两个方法:
* - processBroadcastElement: 处理广播流(限速规则)的每条数据,更新广播状态。
* - processElement: 处理主流(车辆数据)的每条数据,读取广播状态进行判断。
*/
class SpeedCheckBroadcastProcessFunction extends BroadcastProcessFunction[TrafficLog, SpeedLimit, SpeedingViolation] {
override def processBroadcastElement(
limit: SpeedLimit,
ctx: BroadcastProcessFunction[TrafficLog, SpeedLimit, SpeedingViolation]#Context,
out: Collector[SpeedingViolation]
): Unit = {
// 获取广播状态
val state: BroadcastState[String, Int] = ctx.getBroadcastState(SPEED_LIMIT_DESCRIPTOR)
// 以 "monitorId:roadId" 为键,存储限速值
val key = s"${limit.monitorId}:${limit.roadId}"
state.put(key, limit.limitSpeed)
println(s"[Broadcast Updated] $key -> ${limit.limitSpeed} km/h")
}
override def processElement(
log: TrafficLog,
readOnlyCtx: BroadcastProcessFunction[TrafficLog, SpeedLimit, SpeedingViolation]#ReadOnlyContext,
out: Collector[SpeedingViolation]
): Unit = {
// 获取只读的广播状态
val state: ReadOnlyBroadcastState[String, Int] = readOnlyCtx.getBroadcastState(SPEED_LIMIT_DESCRIPTOR)
val key = s"${log.monitorId}:${log.roadId}"
val limitSpeedOpt = Option(state.get(key))
limitSpeedOpt.foreach { limitSpeed =>
// 判断是否超速(这里定义超过限速10%即为超速)
if (log.speed > limitSpeed * 1.1) {
val violation = SpeedingViolation(
carPlate = log.carPlate,
monitorId = log.monitorId,
roadId = log.roadId,
actualSpeed = log.speed,
limitSpeed = limitSpeed,
actionTime = log.actionTime
)
out.collect(violation)
// 这里可以加入更复杂的逻辑,比如同一车辆短时间内多次超速的累加判断
}
}
// 如果找不到对应卡口的限速规则,则忽略此条记录(或按默认规则处理)
}
}
/**
* 自定义Sink,将超速告警写入MySQL。
*/
class JdbcSinkSpeedingViolation extends RichSinkFunction[SpeedingViolation] {
private var connection: Connection = _
private var insertStmt: PreparedStatement = _
override def open(parameters: Configuration): Unit = {
connection = DriverManager.getConnection(
"jdbc:mysql://localhost:3306/traffic_db?useSSL=false&serverTimezone=UTC",
"your_username",
"your_password"
)
val sql =
"""
|INSERT INTO speeding_violation
|(car_plate, monitor_id, road_id, actual_speed, limit_speed, action_time, process_time)
|VALUES (?, ?, ?, ?, ?, ?, ?)
|""".stripMargin
insertStmt = connection.prepareStatement(sql)
}
override def invoke(violation: SpeedingViolation, context: SinkFunction.Context): Unit = {
insertStmt.setString(1, violation.carPlate)
insertStmt.setString(2, violation.monitorId)
insertStmt.setString(3, violation.roadId)
insertStmt.setDouble(4, violation.actualSpeed)
insertStmt.setInt(5, violation.limitSpeed)
insertStmt.setLong(6, violation.actionTime)
insertStmt.setLong(7, violation.processTime)
insertStmt.executeUpdate()
}
override def close(): Unit = {
if (insertStmt != null) insertStmt.close()
if (connection != null) connection.close()
}
}
}
运行这个作业,你会在控制台看到类似这样的输出,同时数据库speeding_violation表中会插入记录:
[Broadcast Updated] 0001:R001 -> 60 km/h
[Broadcast Updated] 0002:R001 -> 80 km/h
...
SpeedingViolation(京B87654,0001,R001,78.5,60,1646123456789,1646123460000)
这表示系统成功检测到一辆车牌为京B87654的车辆在卡口0001(道路R001)超速行驶,实际速度78.5 km/h,超过限速60 km/h的10%。
3. 处理乱序数据与时间窗口聚合
现实世界的数据流很少是严格有序的。网络延迟、设备时钟不同步都可能导致后产生的数据先到达处理系统。Flink通过水印(Watermark) 和事件时间(Event Time) 机制优雅地处理了这个问题。水印可以理解为一种进度指标,它告诉算子“时间戳早于水印的事件应该都已经到达了”,算子可以据此触发窗口计算,而不必无限期等待迟到的数据。
除了超速监控,我们通常还想知道每个卡口的实时平均车速,这能反映路段的拥堵状况。这需要用到Flink的窗口(Window) 计算。我们设计一个滑动窗口:每1分钟统计一次过去5分钟内的数据。
让我们在同一个作业里增加这个功能。我们将创建一个新的流分支来计算平均速度。
// 在 TrafficSpeedMonitor.main 方法中,添加以下代码(在定义 speedingViolationStream 之后)
// 8. 计算每个卡口的平均车速(滑动窗口:窗口大小5分钟,滑动步长1分钟)
val avgSpeedStream: DataStream[MonitorAvgSpeed] = trafficLogStream
.assignTimestampsAndWatermarks(
// 允许数据最大乱序时间为10秒
WatermarkStrategy
.forBoundedOutOfOrderness[TrafficLog](Duration.ofSeconds(10))
.withTimestampAssigner(new SerializableTimestampAssigner[TrafficLog] {
override def extractTimestamp(element: TrafficLog, recordTimestamp: Long): Long = {
element.actionTime // 使用过车时间作为事件时间
}
})
)
.keyBy(log => log.monitorId) // 按卡口ID分组
.window(
SlidingEventTimeWindows.of(Time.minutes(5), Time.minutes(1))
)
.aggregate(new AvgSpeedAggregator, new AvgSpeedWindowProcessor)
.name("avg-speed-calculator")
// 定义平均速度的结果类
case class MonitorAvgSpeed(
windowStart: Long,
windowEnd: Long,
monitorId: String,
avgSpeed: Double,
vehicleCount: Int
)
// 增量聚合函数:累加车辆数和速度总和
class AvgSpeedAggregator extends AggregateFunction[TrafficLog, (Int, Double), (Int, Double)] {
override def createAccumulator(): (Int, Double) = (0, 0.0)
override def add(value: TrafficLog, accumulator: (Int, Double)): (Int, Double) = {
(accumulator._1 + 1, accumulator._2 + value.speed)
}
override def getResult(accumulator: (Int, Double)): (Int, Double) = accumulator
override def merge(a: (Int, Double), b: (Int, Double)): (Int, Double) = {
(a._1 + b._1, a._2 + b._2)
}
}
// 窗口处理函数:计算最终的平均值
class AvgSpeedWindowProcessor extends ProcessWindowFunction[(Int, Double), MonitorAvgSpeed, String, TimeWindow] {
override def process(
key: String,
context: Context,
elements: Iterable[(Int, Double)],
out: Collector[MonitorAvgSpeed]
): Unit = {
val (count, totalSpeed) = elements.head // 因为增量聚合后,Iterable里只有一个元素
val avgSpeed = if (count > 0) totalSpeed / count else 0.0
out.collect(MonitorAvgSpeed(
context.window.getStart,
context.window.getEnd,
key,
avgSpeed,
count
))
}
}
// 9. 将平均速度结果也写入另一个MySQL表或打印出来
avgSpeedStream.addSink(new JdbcSinkAvgSpeed).name("jdbc-sink-avg-speed")
avgSpeedStream.print().name("console-sink-avg-speed")
// 对应的JdbcSinkAvgSpeed类(结构与JdbcSinkSpeedingViolation类似,略)
这个设计有几个关键点:
forBoundedOutOfOrderness(Duration.ofSeconds(10)):设置了10秒的最大乱序容忍度,即允许数据迟到10秒。这是一个需要根据数据源特性调整的经验值。- 滑动窗口:
SlidingEventTimeWindows.of(Time.minutes(5), Time.minutes(1))意味着每1分钟输出一次过去5分钟的计算结果。这比滚动窗口(Tumbling Window)能提供更平滑、更实时的趋势视图。 - 增量聚合:
AggregateFunction在数据到达时即进行累加(add方法),而不是等到窗口触发时才遍历所有原始数据,这大大降低了状态存储的压力和计算开销。
4. 部署优化与生产环境考量
代码在本地跑通只是第一步。要将其部署到生产环境,成为一个健壮的城市交通监控平台组件,我们还需要考虑更多。下面这个表格对比了开发原型与生产系统需要关注的不同维度:
| 考量维度 | 开发/原型阶段 | 生产环境建议 |
|---|---|---|
| 数据源 | 模拟数据源 (SimulatedTrafficSource) | 接入真实Kafka集群,主题分区策略需根据卡口ID或区域ID设计,以并行消费。 |
| 状态后端 | 默认的MemoryStateBackend(易失) | RocksDBStateBackend,将状态持久化到本地磁盘或HDFS,支持大状态和异步快照。 |
| 检查点与容错 | 未开启或间隔很长 | 开启Checkpointing(间隔如1分钟),并配置外部持久化存储(如HDFS、S3)保存状态快照。作业失败后可从最近检查点恢复。 |
| 水位线生成 | 简单的周期性生成 | 根据数据流特征定制,对于稀疏流可使用WatermarkStrategy.forMonotonousTimestamps,对于乱序严重的流可使用forBoundedOutOfOrderness并合理设置延迟。 |
| 维度表更新 | 定时全量拉取(广播状态) | 对于更新频繁的维度表(如限速规则),可考虑使用CDC(Change Data Capture) 工具(如Debezium)捕获MySQL的binlog,作为另一条流与主流连接,实现更实时的规则更新。 |
| 结果输出 | 单点MySQL | 根据下游需求多样化:实时告警可推送至Kafka(供其他系统订阅)、Redis(供前端实时展示);批量分析可写入ClickHouse或HBase;持久化存储仍可保留MySQL。 |
| 监控与运维 | 控制台打印日志 | 集成Metrics系统(如Prometheus + Grafana),监控作业吞吐量、延迟、背压、Checkpoint时长等关键指标。配置报警规则。 |
| 资源与并行度 | 并行度=1 | 根据数据量和集群资源设置合适的并行度。KeyBy后的算子并行度会受Key分布影响,需避免数据倾斜。 |
| 代码结构 | 单作业,所有逻辑在一起 | 遵循单一职责原则,可将超速检测、平均速度计算、拥堵排名等拆分为独立的Flink作业,通过Kafka主题连接,形成流处理管道,提高可维护性和可扩展性。 |
提示:在设置Kafka消费者时,建议将
auto.offset.reset策略设置为latest(从最新消费)或earliest(从最早消费),并在作业中启用检查点(Checkpoint) 来提交偏移量,这样可以保证在作业重启时能从故障点继续消费,避免数据丢失或重复。
对于维度表(如限速规则)的更新,除了我们使用的广播状态模式,还有一种更高级的模式是时态表连接(Temporal Table Join)。它特别适用于维度表有版本变更历史,且需要将数据流与特定时间点的维度表版本进行关联的场景。例如,我们需要知道车辆超速时,当时生效的限速规则是哪一条。这需要维度表包含时间版本字段(如start_time, end_time)。虽然实现稍复杂,但能提供更精确的历史关联。
最后,分享一个我在实际项目中遇到的坑:Key的设计。最初我用monitorId作为广播状态的Key,后来发现同一monitorId可能对应不同道路(比如一个十字路口有四个方向),导致规则错乱。所以最终使用了monitorId:roadId这个组合键。在设计Key时,一定要确保其能唯一标识一条业务规则。另一个经验是,广播状态不宜过大,因为它会复制到每个并行子任务中。如果维度表数据量巨大(例如百万级),广播状态模式可能不是最佳选择,需要考虑其他方案如Async I/O。
更多推荐
所有评论(0)