共计 3352 个字符,预计需要花费 9 分钟才能阅读完成。
背景痛点:为什么 AIS 数据这么难处理?
AIS(Automatic Identification System)数据是船舶自动发送的航行信息,包含位置、航速、航向等关键数据。这类数据有几个显著特点:

- 数据量大:全球有超过 10 万艘商船,每艘船每 2 -10 秒就会发送一条 AIS 消息
- 格式复杂:原始数据包含 NMEA 0183 二进制格式、CSV、JSON 等多种形式
- 时效性高:碰撞预警等场景需要秒级延迟
传统单机处理方式(如用 Python+pandas)在处理 TB 级数据时面临三大难题:
- 内存不足导致 OOM(内存溢出)
- 处理速度跟不上数据生成速度
- 难以实现复杂的实时分析逻辑
技术选型:Spark+Flink 混合架构
经过对比测试,我们最终选择 Spark+Flink 的混合方案:
| 技术 | 适用场景 | AIS 数据处理中的应用 |
|---|---|---|
| Spark | 批量 ETL | 历史数据清洗、格式标准化 |
| Flink | 流处理 | 实时轨迹分析、异常检测 |
混合架构优势:
- 资源利用率高:白天用 Flink 处理实时数据,夜间用 Spark 跑批量作业
- 开发效率高:两者都支持 SQL 接口,代码可复用
- 运维简单:YARN/K8s 上可统一管理
核心实现:从原始数据到业务洞察
阶段一:Spark 数据清洗(Scala 示例)
// 读取原始 NMEA 格式数据
val rawDF = spark.read.format("text").load("s3://ais-bucket/raw/*.log")
// 定义 UDF 解析 NMEA 报文
val parseNmea = udf {(line: String) =>
// $GPAAM,A,A,0.10,N,WPTNME*32
try {val parts = line.split(",")
Map("msg_type" -> parts(0).substring(1),
"mmsi" -> parts(1),
"lat" -> parts(2).toDouble,
"lon" -> parts(3).toDouble
)
} catch {case _: Exception => null // 脏数据处理}
}
// 数据清洗管道
val cleanDF = rawDF
.filter(col("value").startsWith("$")) // 过滤无效行
.withColumn("parsed", parseNmea(col("value")))
.select("parsed.*")
.na.drop() // 删除解析失败记录
// 保存为 Parquet 格式(按日期分区)cleanDF.write.partitionBy("date").parquet("s3://ais-bucket/cleaned/")
关键点说明:
- 使用 UDF 处理复杂解析逻辑
- 采用 filter+na.drop 组合拳清理脏数据
- Parquet 分区提升查询性能
阶段二:Flink 实时分析(Java 示例)
// 定义 Kafka 数据源
KafkaSource<String> source = KafkaSource.<String>builder()
.setBootstrapServers("kafka:9092")
.setTopics("ais-stream")
.setDeserializer(new SimpleStringSchema())
.build();
// 构建处理管道
DataStream<AISEvent> events = env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Source")
.map(new JSONParser()) // 解析 JSON
.keyBy(event -> event.getMmsi()) // 按船舶 ID 分组
.window(TumblingEventTimeWindows.of(Time.minutes(5))) // 5 分钟窗口
.process(new ShipTrajectoryAnalyzer());
// 自定义窗口函数
public static class ShipTrajectoryAnalyzer
extends ProcessWindowFunction<AISEvent, Alert, String, TimeWindow> {
@Override
public void process(
String mmsi,
Context context,
Iterable<AISEvent> events,
Collector<Alert> out) {
// 状态管理:通过 ListState 保存历史位置
ListState<Position> positionState = context.globalState()
.getListState(new ListStateDescriptor<>("positions", Position.class));
// 计算平均航速、方向变化率等指标
// 触发碰撞预警规则
if (detectCollisionRisk(events)) {out.collect(new Alert(mmsi, "COLLISION_RISK"));
}
}
}
实时处理要点:
- 使用 Keyed Window 确保同一船舶的数据进入相同计算节点
- 通过 State API 管理船舶历史轨迹
- 自定义 ProcessFunction 实现复杂业务逻辑
性能优化实战经验
存储优化方案
- Parquet 压缩:采用 SNAPPY 压缩(CPU 开销与压缩率的平衡)
df.write.option("compression", "snappy").parquet("output_path") - 分区策略 :按
年 / 月 / 日三级分区,避免单个分区过大 - 小文件合并:每小时执行一次小文件合并
spark.read.parquet("path") .repartition(10) // 按数据量调整 .write.parquet("new_path")
Flink 资源配置黄金法则
- 并行度设置:
- Kafka 分区数 = Flink 并行度
-
每个 TaskManager 的 slot 数 = CPU 核数 – 1(留出系统资源)
-
内存分配(以 4G 容器为例):
taskmanager.memory.process.size: 4096m taskmanager.memory.task.heap.size: 2048m # 50% 给堆内存 taskmanager.memory.managed.size: 1024m # 25% 给 Flink 管理内存 -
检查点调优:
env.enableCheckpointing(30000); // 30 秒一次 env.getCheckpointConfig().setMinPauseBetweenCheckpoints(10000);
避坑指南:血泪教训总结
时间戳处理三原则
- 统一时区:全部转换为 UTC 时间存储
- 事件时间:必须配置 Watermark(允许延迟)
WatermarkStrategy<Event> strategy = WatermarkStrategy .<Event>forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, ts) -> event.getTimestamp()); - 批流一致:Spark 和 Flink 使用相同的时间解析逻辑
状态后端选型
| 类型 | 特点 | 适用场景 |
|---|---|---|
| MemoryStateBackend | 本地内存,宕机丢失 | 测试环境 |
| FsStateBackend | 内存 + 文件系统 | 生产环境通用 |
| RocksDBStateBackend | 磁盘存储,支持大状态 | 超大规模状态 |
推荐配置:
env.setStateBackend(new RocksDBStateBackend("hdfs://checkpoints", true));
方案总结与扩展思考
这套方案在某海运监控平台的实际表现:
- 吞吐量:单集群处理 20 万条 / 秒
- 延迟:从数据产生到预警平均 1.2 秒
- 存储效率:压缩比达到 1:8
扩展建议:
- 设备维保预测:结合发动机转速等数据预测故障
- 港口调度优化:分析靠泊时间规律
- 碳排放计算:根据航速估算燃油消耗
处理 AIS 数据的核心思想——” 分而治之 ” 同样适用于其他物联网场景:
- 车联网:用相同架构处理 GPS 数据
- 工业传感器:振动数据实时分析
- 智能电表:用电量异常检测
关键是要根据业务特点调整窗口策略和状态管理方式。
正文完
