diff --git a/ingest-api/src/main/java/com/lingniu/ingest/api/event/VehicleEvent.java b/ingest-api/src/main/java/com/lingniu/ingest/api/event/VehicleEvent.java index 1044b8a4..58fb23a3 100644 --- a/ingest-api/src/main/java/com/lingniu/ingest/api/event/VehicleEvent.java +++ b/ingest-api/src/main/java/com/lingniu/ingest/api/event/VehicleEvent.java @@ -26,7 +26,8 @@ public sealed interface VehicleEvent VehicleEvent.Logout, VehicleEvent.Heartbeat, VehicleEvent.MediaMeta, - VehicleEvent.Passthrough { + VehicleEvent.Passthrough, + VehicleEvent.RawArchive { String eventId(); String vin(); @@ -138,4 +139,28 @@ public sealed interface VehicleEvent int passthroughType, byte[] data ) implements VehicleEvent {} + + /** + * 原始报文冷存事件:每条成功解码的入站帧由 Dispatcher 产出一条,携带原始字节交 + * {@code ArchiveEventSink} 写入 ArchiveStore。Kafka sink 默认不处理本类型 + * (见 {@code KafkaEventSink.accepts})。 + * + *

key 组装建议:{@code yyyy/MM/dd///.bin},具体由 sink 实现决定。 + * + * @param command 协议主命令码(如 32960 0x02/0x03 等),冗余在事件里便于按命令分片归档 + * @param infoType 协议子类型(可为 0),同上 + * @param rawBytes 原始入站字节,不可为 null + */ + record RawArchive( + String eventId, + String vin, + ProtocolId source, + Instant eventTime, + Instant ingestTime, + String traceId, + Map metadata, + int command, + int infoType, + byte[] rawBytes + ) implements VehicleEvent {} } diff --git a/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/EnvelopeMapper.java b/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/EnvelopeMapper.java index 3eb4a645..e2684279 100644 --- a/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/EnvelopeMapper.java +++ b/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/EnvelopeMapper.java @@ -57,6 +57,15 @@ public final class EnvelopeMapper { .setData(com.google.protobuf.ByteString.copyFrom( p.data() == null ? new byte[0] : p.data())) .build()); + case VehicleEvent.RawArchive ra -> { + // RawArchive 原则上由 ArchiveEventSink 独占消费,本分支仅在有人显式把它送到 + // Kafka 时生效:不设 payload oneof,只填 raw_archive ref 里的 size_bytes, + // 下游如需真正拉取字节要依赖 archive 写成功后回填的 URI(后续迭代)。 + int size = ra.rawBytes() == null ? 0 : ra.rawBytes().length; + b.setRawArchive(com.lingniu.ingest.sink.mq.proto.RawArchiveRef.newBuilder() + .setSizeBytes(size) + .build()); + } } return b.build(); } diff --git a/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEventSink.java b/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEventSink.java index f2dd4540..4a05df81 100644 --- a/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEventSink.java +++ b/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEventSink.java @@ -42,6 +42,16 @@ public final class KafkaEventSink implements EventSink, AutoCloseable { return "kafka"; } + /** + * 拒绝 {@link VehicleEvent.RawArchive} —— 原始报文体积大、属于冷存域,本期由 + * {@code ArchiveEventSink} 独占处理。未来若需要将 raw archive URI 回填到 Kafka 的 + * {@code vehicle.raw.archive} topic,再放开此过滤并在 Envelope 里填 uri/checksum。 + */ + @Override + public boolean accepts(VehicleEvent event) { + return !(event instanceof VehicleEvent.RawArchive); + } + @Override public CompletableFuture publish(VehicleEvent event) { CompletableFuture cf = new CompletableFuture<>(); diff --git a/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/TopicRouter.java b/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/TopicRouter.java index 2c22b7c9..b4a7cbc0 100644 --- a/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/TopicRouter.java +++ b/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/TopicRouter.java @@ -23,6 +23,7 @@ public final class TopicRouter { VehicleEvent.Heartbeat _ -> topics.getSession(); case VehicleEvent.MediaMeta _ -> topics.getMediaMeta(); case VehicleEvent.Passthrough _ -> topics.getAlarm(); + case VehicleEvent.RawArchive _ -> topics.getRawArchive(); }; } }