From 98b37e71e926a19498bff02769e30f9d6a040c7f Mon Sep 17 00:00:00 2001 From: lingniu Date: Mon, 29 Jun 2026 15:56:41 +0800 Subject: [PATCH] fix: keep jt808 raw history consumer current --- ...icleHistoryKafkaConsumerConfiguration.java | 12 ++- .../src/main/resources/application.yml | 11 ++- .../VehicleHistoryAppDefaultsTest.java | 20 +++-- .../consumer/EnvelopeConsumerProcessor.java | 33 ++++++- .../ingest/api/consumer/EnvelopeIngestor.java | 14 +++ .../EventHistoryEnvelopeIngestor.java | 88 +++++++++++++------ .../EventHistoryEnvelopeIngestorTest.java | 73 +++++++++++++++ .../sink/mq/KafkaEnvelopeConsumerFactory.java | 1 + .../sink/mq/KafkaEnvelopeConsumerWorker.java | 17 ++-- .../ingest/sink/mq/SinkMqProperties.java | 3 + .../mq/KafkaEnvelopeConsumerFactoryTest.java | 1 + .../TdengineJdbcHistoryWriter.java | 36 +++++--- .../TdengineJdbcHistoryWriterTest.java | 19 ++++ 13 files changed, 274 insertions(+), 54 deletions(-) diff --git a/modules/apps/vehicle-history-app/src/main/java/com/lingniu/ingest/historyapp/VehicleHistoryKafkaConsumerConfiguration.java b/modules/apps/vehicle-history-app/src/main/java/com/lingniu/ingest/historyapp/VehicleHistoryKafkaConsumerConfiguration.java index 754747f7..69fd7f1d 100644 --- a/modules/apps/vehicle-history-app/src/main/java/com/lingniu/ingest/historyapp/VehicleHistoryKafkaConsumerConfiguration.java +++ b/modules/apps/vehicle-history-app/src/main/java/com/lingniu/ingest/historyapp/VehicleHistoryKafkaConsumerConfiguration.java @@ -45,14 +45,24 @@ public class VehicleHistoryKafkaConsumerConfiguration { return new EnvelopeConsumerProcessor("event-history", ingestor, deadLetterSink); } + @Bean + @ConditionalOnMissingBean(name = "eventHistoryRawEnvelopeConsumerProcessor") + public EnvelopeConsumerProcessor eventHistoryRawEnvelopeConsumerProcessor(EventHistoryEnvelopeIngestor ingestor, + EnvelopeDeadLetterSink deadLetterSink) { + return new EnvelopeConsumerProcessor("event-history-raw", ingestor, deadLetterSink); + } + @Bean @ConditionalOnMissingBean @ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq.consumer", name = "enabled", havingValue = "true") public KafkaEnvelopeConsumerRunner vehicleHistoryKafkaEnvelopeConsumerRunner( @Qualifier("eventHistoryEnvelopeConsumerProcessor") EnvelopeConsumerProcessor processor, + @Qualifier("eventHistoryRawEnvelopeConsumerProcessor") EnvelopeConsumerProcessor rawProcessor, SinkMqProperties props) { List workers = new KafkaEnvelopeConsumerFactory().createWorkers( - Map.of("eventHistoryEnvelopeConsumerProcessor", processor), + Map.of( + "eventHistoryEnvelopeConsumerProcessor", processor, + "eventHistoryRawEnvelopeConsumerProcessor", rawProcessor), props); if (workers.isEmpty()) { throw new IllegalStateException("no vehicle history kafka consumer workers created; check consumer bindings"); diff --git a/modules/apps/vehicle-history-app/src/main/resources/application.yml b/modules/apps/vehicle-history-app/src/main/resources/application.yml index b6ec8d2e..867fc42e 100644 --- a/modules/apps/vehicle-history-app/src/main/resources/application.yml +++ b/modules/apps/vehicle-history-app/src/main/resources/application.yml @@ -46,15 +46,20 @@ lingniu: enabled: ${KAFKA_CONSUMER_ENABLED:true} client-id-prefix: ${KAFKA_CONSUMER_CLIENT_ID_PREFIX:vehicle-history} auto-offset-reset: ${KAFKA_CONSUMER_AUTO_OFFSET_RESET:earliest} - max-poll-records: ${KAFKA_CONSUMER_MAX_POLL_RECORDS:500} + max-poll-records: ${KAFKA_CONSUMER_MAX_POLL_RECORDS:100} + max-poll-interval-millis: ${KAFKA_CONSUMER_MAX_POLL_INTERVAL_MS:1800000} bindings: eventHistoryEnvelopeConsumerProcessor: enabled: true - group-id: ${KAFKA_GROUP_HISTORY:vehicle-history} + group-id: ${KAFKA_GROUP_HISTORY_EVENT:${KAFKA_GROUP_HISTORY:vehicle-history}-event} topics: - ${KAFKA_TOPIC_GB32960_EVENT:vehicle.event.gb32960.v1} - - ${KAFKA_TOPIC_GB32960_RAW:vehicle.raw.gb32960.v1} - ${KAFKA_TOPIC_JT808_EVENT:vehicle.event.jt808.v1} + eventHistoryRawEnvelopeConsumerProcessor: + enabled: true + group-id: ${KAFKA_GROUP_HISTORY_RAW:${KAFKA_GROUP_HISTORY:vehicle-history}} + topics: + - ${KAFKA_TOPIC_GB32960_RAW:vehicle.raw.gb32960.v1} - ${KAFKA_TOPIC_JT808_RAW:vehicle.raw.jt808.v1} archive: enabled: ${SINK_ARCHIVE_ENABLED:true} diff --git a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppDefaultsTest.java b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppDefaultsTest.java index 92aa42ab..0b2b4b01 100644 --- a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppDefaultsTest.java +++ b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppDefaultsTest.java @@ -40,23 +40,33 @@ class VehicleHistoryAppDefaultsTest { .containsEntry("lingniu.ingest.vehicle-state.enabled", false) .containsEntry("lingniu.ingest.vehicle-stat.enabled", false) .containsEntry("lingniu.ingest.sink.mq.consumer.enabled", "${KAFKA_CONSUMER_ENABLED:true}") + .containsEntry("lingniu.ingest.sink.mq.consumer.max-poll-records", "${KAFKA_CONSUMER_MAX_POLL_RECORDS:100}") + .containsEntry( + "lingniu.ingest.sink.mq.consumer.max-poll-interval-millis", + "${KAFKA_CONSUMER_MAX_POLL_INTERVAL_MS:1800000}") .containsEntry( "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryEnvelopeConsumerProcessor.enabled", true) .containsEntry( "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryEnvelopeConsumerProcessor.group-id", - "${KAFKA_GROUP_HISTORY:vehicle-history}") + "${KAFKA_GROUP_HISTORY_EVENT:${KAFKA_GROUP_HISTORY:vehicle-history}-event}") .containsEntry( "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryEnvelopeConsumerProcessor.topics[0]", "${KAFKA_TOPIC_GB32960_EVENT:vehicle.event.gb32960.v1}") .containsEntry( "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryEnvelopeConsumerProcessor.topics[1]", - "${KAFKA_TOPIC_GB32960_RAW:vehicle.raw.gb32960.v1}") - .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryEnvelopeConsumerProcessor.topics[2]", "${KAFKA_TOPIC_JT808_EVENT:vehicle.event.jt808.v1}") .containsEntry( - "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryEnvelopeConsumerProcessor.topics[3]", + "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryRawEnvelopeConsumerProcessor.enabled", + true) + .containsEntry( + "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryRawEnvelopeConsumerProcessor.group-id", + "${KAFKA_GROUP_HISTORY_RAW:${KAFKA_GROUP_HISTORY:vehicle-history}}") + .containsEntry( + "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryRawEnvelopeConsumerProcessor.topics[0]", + "${KAFKA_TOPIC_GB32960_RAW:vehicle.raw.gb32960.v1}") + .containsEntry( + "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryRawEnvelopeConsumerProcessor.topics[1]", "${KAFKA_TOPIC_JT808_RAW:vehicle.raw.jt808.v1}"); } diff --git a/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/consumer/EnvelopeConsumerProcessor.java b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/consumer/EnvelopeConsumerProcessor.java index 68b4041b..7359179f 100644 --- a/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/consumer/EnvelopeConsumerProcessor.java +++ b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/consumer/EnvelopeConsumerProcessor.java @@ -1,7 +1,9 @@ package com.lingniu.ingest.api.consumer; import java.time.Instant; +import java.util.ArrayList; import java.util.EnumSet; +import java.util.List; import java.util.Set; public final class EnvelopeConsumerProcessor { @@ -32,13 +34,40 @@ public final class EnvelopeConsumerProcessor { if (record == null) { throw new IllegalArgumentException("record must not be null"); } - EnvelopeIngestResult result = ingestor.tryIngest(record.payload()); + EnvelopeIngestResult result = processAll(List.of(record)).getFirst(); + return result; + } + + public List processAll(List records) { + if (records == null) { + throw new IllegalArgumentException("records must not be null"); + } + if (records.isEmpty()) { + return List.of(); + } + List payloads = new ArrayList<>(records.size()); + for (EnvelopeConsumerRecord record : records) { + if (record == null) { + throw new IllegalArgumentException("record must not be null"); + } + payloads.add(record.payload()); + } + List results = ingestor.tryIngestAll(payloads); + if (results.size() != records.size()) { + throw new IllegalStateException("ingestor result count does not match record count"); + } + for (int i = 0; i < records.size(); i++) { + publishDeadLetterIfNeeded(records.get(i), results.get(i)); + } + return results; + } + + private void publishDeadLetterIfNeeded(EnvelopeConsumerRecord record, EnvelopeIngestResult result) { // 派生消费者不在这里抛出业务异常给 Kafka worker;坏消息统一进入 DLQ, // worker 看到 process 正常返回后才能提交 offset,避免同一坏消息无限阻塞消费组。 if (DEAD_LETTER_STATUSES.contains(result.status())) { deadLetterSink.publish(toDeadLetter(record, result)); } - return result; } private EnvelopeDeadLetterRecord toDeadLetter(EnvelopeConsumerRecord record, EnvelopeIngestResult result) { diff --git a/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/consumer/EnvelopeIngestor.java b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/consumer/EnvelopeIngestor.java index 471dae68..2d200465 100644 --- a/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/consumer/EnvelopeIngestor.java +++ b/modules/core/ingest-api/src/main/java/com/lingniu/ingest/api/consumer/EnvelopeIngestor.java @@ -1,6 +1,20 @@ package com.lingniu.ingest.api.consumer; +import java.util.ArrayList; +import java.util.List; + @FunctionalInterface public interface EnvelopeIngestor { EnvelopeIngestResult tryIngest(byte[] kafkaValue); + + default List tryIngestAll(List kafkaValues) { + if (kafkaValues == null || kafkaValues.isEmpty()) { + return List.of(); + } + List results = new ArrayList<>(kafkaValues.size()); + for (byte[] kafkaValue : kafkaValues) { + results.add(tryIngest(kafkaValue)); + } + return results; + } } diff --git a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestor.java b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestor.java index 5a3acc2b..5c1ba757 100644 --- a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestor.java +++ b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestor.java @@ -10,6 +10,7 @@ import com.lingniu.ingest.tdenginehistory.TdengineEnvelopeRows; import com.lingniu.ingest.tdenginehistory.TdengineHistoryWriter; import java.io.IOException; +import java.util.ArrayList; import java.util.List; /** @@ -47,30 +48,57 @@ public final class EventHistoryEnvelopeIngestor implements EnvelopeIngestor { VehicleEnvelope envelope = parse(kafkaValue); EventFileRecord record = mapper.toRecord(envelope); store.append(record); - writeTdengineFacts(envelope); + writeTdengineFacts(List.of(envelope)); } @Override public EnvelopeIngestResult tryIngest(byte[] kafkaValue) { - VehicleEnvelope envelope = null; - try { - envelope = parse(kafkaValue); - EventFileRecord record = mapper.toRecord(envelope); - store.append(record); - writeTdengineFacts(envelope); - return EnvelopeIngestResult.stored(record.eventId(), record.vin()); - } catch (IllegalArgumentException ex) { - // 协议不可解析或缺少必要字段属于不可重试问题,避免 Kafka consumer 无限重放。 - return envelope == null - ? EnvelopeIngestResult.invalid(ex.getMessage()) - : EnvelopeIngestResult.skipped(envelope.getEventId(), envelope.getVin(), ex.getMessage()); - } catch (IOException ex) { - // 文件库写入失败可能是临时 I/O 问题,交给上层 worker 按失败处理。 - return EnvelopeIngestResult.failed( - envelope == null ? "" : envelope.getEventId(), - envelope == null ? "" : envelope.getVin(), - ex.getMessage()); + return tryIngestAll(List.of(kafkaValue)).getFirst(); + } + + @Override + public List tryIngestAll(List kafkaValues) { + if (kafkaValues == null || kafkaValues.isEmpty()) { + return List.of(); } + List results = new ArrayList<>(kafkaValues.size()); + List valid = new ArrayList<>(kafkaValues.size()); + for (byte[] kafkaValue : kafkaValues) { + VehicleEnvelope envelope = null; + try { + envelope = parse(kafkaValue); + EventFileRecord record = mapper.toRecord(envelope); + results.add(null); + valid.add(new BatchEntry(results.size() - 1, envelope, record)); + } catch (IllegalArgumentException ex) { + results.add(envelope == null + ? EnvelopeIngestResult.invalid(ex.getMessage()) + : EnvelopeIngestResult.skipped(envelope.getEventId(), envelope.getVin(), ex.getMessage())); + } + } + if (valid.isEmpty()) { + return List.copyOf(results); + } + try { + store.appendAll(valid.stream().map(BatchEntry::record).toList()); + writeTdengineFacts(valid.stream().map(BatchEntry::envelope).toList()); + for (BatchEntry entry : valid) { + results.set(entry.index(), EnvelopeIngestResult.stored( + entry.record().eventId(), + entry.record().vin())); + } + } catch (IOException ex) { + for (BatchEntry entry : valid) { + results.set(entry.index(), EnvelopeIngestResult.failed( + entry.envelope().getEventId(), + entry.envelope().getVin(), + ex.getMessage())); + } + } + return List.copyOf(results); + } + + private record BatchEntry(int index, VehicleEnvelope envelope, EventFileRecord record) { } private static VehicleEnvelope parse(byte[] kafkaValue) { @@ -84,19 +112,25 @@ public final class EventHistoryEnvelopeIngestor implements EnvelopeIngestor { } } - private void writeTdengineFacts(VehicleEnvelope envelope) throws IOException { + private void writeTdengineFacts(List envelopes) throws IOException { if (tdengineWriter == null) { return; } - var rawFrame = TdengineEnvelopeRows.rawFrame(envelope); - if (rawFrame.isPresent()) { - tdengineWriter.appendRawFrames(List.of(rawFrame.get())); + var rawFrames = envelopes.stream() + .flatMap(envelope -> TdengineEnvelopeRows.rawFrame(envelope).stream()) + .toList(); + if (!rawFrames.isEmpty()) { + tdengineWriter.appendRawFrames(rawFrames); } - var location = TdengineEnvelopeRows.location(envelope); - if (location.isPresent()) { - tdengineWriter.appendLocations(List.of(location.get())); + var locations = envelopes.stream() + .flatMap(envelope -> TdengineEnvelopeRows.location(envelope).stream()) + .toList(); + if (!locations.isEmpty()) { + tdengineWriter.appendLocations(locations); } - var telemetryFields = TdengineEnvelopeRows.telemetryFields(envelope); + var telemetryFields = envelopes.stream() + .flatMap(envelope -> TdengineEnvelopeRows.telemetryFields(envelope).stream()) + .toList(); if (!telemetryFields.isEmpty()) { tdengineWriter.appendTelemetryFields(telemetryFields); } diff --git a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestorTest.java b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestorTest.java index cc7d8526..f91e7f99 100644 --- a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestorTest.java +++ b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/EventHistoryEnvelopeIngestorTest.java @@ -186,6 +186,38 @@ class EventHistoryEnvelopeIngestorTest { assertThat(tdengineWriter.telemetryFields.getFirst().valueDouble()).isEqualTo(42.5); } + @Test + void tryIngestAllBatchesFileStoreAndTdengineWrites() throws Exception { + CapturingStore store = new CapturingStore(); + CapturingTdengineWriter tdengineWriter = new CapturingTdengineWriter(); + EventHistoryEnvelopeIngestor ingestor = new EventHistoryEnvelopeIngestor( + store, new TelemetryEnvelopeRecordMapper(), tdengineWriter); + + VehicleEnvelope first = jt808LocationEnvelope("jt808-location-1", "frame-jt808-1", "013800000001"); + VehicleEnvelope second = jt808LocationEnvelope("jt808-location-2", "frame-jt808-2", "013800000002"); + + List results = ingestor.tryIngestAll(List.of( + first.toByteArray(), + new byte[]{0x01, 0x02}, + second.toByteArray())); + + assertThat(results).extracting(EnvelopeIngestResult::status) + .containsExactly( + EnvelopeIngestResult.Status.STORED, + EnvelopeIngestResult.Status.INVALID_ENVELOPE, + EnvelopeIngestResult.Status.STORED); + assertThat(store.appendAllCalls).isEqualTo(1); + assertThat(store.records).extracting(EventFileRecord::eventId) + .containsExactly("jt808-location-1", "jt808-location-2"); + assertThat(tdengineWriter.rawFrames).extracting(TdengineRawFrameRow::frameId) + .containsExactly("frame-jt808-1", "frame-jt808-2"); + assertThat(tdengineWriter.locations).extracting(TdengineLocationRow::factId) + .containsExactly("jt808-location-1", "jt808-location-2"); + assertThat(tdengineWriter.telemetryFields) + .extracting(TdengineTelemetryFieldRow::fieldKey) + .containsExactly("location.speedKmh", "location.speedKmh"); + } + private static VehicleEnvelope envelope(String eventId) { return VehicleEnvelope.newBuilder() .setEventId(eventId) @@ -204,11 +236,52 @@ class EventHistoryEnvelopeIngestorTest { .build(); } + private static VehicleEnvelope jt808LocationEnvelope(String eventId, String frameId, String phone) { + return VehicleEnvelope.newBuilder() + .setSchemaVersion("1.0") + .setEventId(eventId) + .setVin("VINJT808001") + .setSource("JT808") + .setEventTimeMs(1_782_112_400_000L) + .setIngestTimeMs(1_782_112_401_000L) + .putMetadata("vehicle_key", "jt808:" + phone) + .putMetadata("frame_id", frameId) + .setRawArchive(RawArchiveRef.newBuilder() + .setUri("archive://jt808/2026/06/29/" + frameId + ".bin") + .setSizeBytes(68)) + .setRawFrameFact(RawFrameFactPayload.newBuilder() + .setFrameId(frameId) + .setVehicleKey("jt808:" + phone) + .setVin("VINJT808001") + .setPhone(phone) + .setMessageId(0x0200) + .setRawUri("archive://jt808/2026/06/29/" + frameId + ".bin") + .setRawSizeBytes(68) + .setParseStatus(ParseStatusProto.PARSE_STATUS_SUCCEEDED)) + .setLocation(LocationPayload.newBuilder() + .setLongitude(113.12) + .setLatitude(23.45) + .setSpeedKmh(42.5)) + .setTelemetrySnapshot(TelemetrySnapshot.newBuilder() + .setEventType("LOCATION") + .setRawArchiveUri("archive://jt808/2026/06/29/" + frameId + ".bin") + .addFields(TelemetryField.newBuilder() + .setKey("location.speedKmh") + .setValueType("DOUBLE") + .setValue("42.5") + .setUnit("km/h") + .setQuality("GOOD") + .setSourcePath("JT808.0x0200.speed"))) + .build(); + } + private static final class CapturingStore implements EventFileStore { private final List records = new ArrayList<>(); + private int appendAllCalls; @Override public void appendAll(List records) { + appendAllCalls++; this.records.addAll(records); } diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerFactory.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerFactory.java index d85e37e0..e0123cda 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerFactory.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerFactory.java @@ -117,6 +117,7 @@ public final class KafkaEnvelopeConsumerFactory { p.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); p.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, consumer.getAutoOffsetReset()); p.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, consumer.getMaxPollRecords()); + p.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, consumer.getMaxPollIntervalMillis()); return p; } diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerWorker.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerWorker.java index 874c1ca1..9e8fabe0 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerWorker.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerWorker.java @@ -7,7 +7,10 @@ import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import java.time.Duration; +import java.util.ArrayList; import java.util.Collection; +import java.util.LinkedHashMap; +import java.util.List; import java.util.Map; public final class KafkaEnvelopeConsumerWorker implements AutoCloseable { @@ -33,22 +36,26 @@ public final class KafkaEnvelopeConsumerWorker implements AutoCloseable { public int pollOnce(Duration timeout) { ConsumerRecords records = consumer.poll(timeout == null ? Duration.ZERO : timeout); - int processed = 0; + Map> batches = new LinkedHashMap<>(); for (ConsumerRecord record : records) { EnvelopeConsumerProcessor processor = processorsByTopic.get(record.topic()); if (processor == null) { // worker 可能订阅多个 topic;没有显式绑定处理器的 topic 不参与提交语义。 continue; } - // EnvelopeConsumerProcessor 内部会把解析或业务错误转成 DLQ 记录, - // 这里保持 Kafka worker 的职责单一:轮询、分发、成功后提交 offset。 - processor.process(new EnvelopeConsumerRecord( + batches.computeIfAbsent(processor, ignored -> new ArrayList<>()).add(new EnvelopeConsumerRecord( record.topic(), record.partition(), record.offset(), record.key(), record.value())); - processed++; + } + int processed = 0; + for (Map.Entry> entry : batches.entrySet()) { + // EnvelopeConsumerProcessor 内部会把解析或业务错误转成 DLQ 记录, + // 这里保持 Kafka worker 的职责单一:轮询、批量分发、成功后提交 offset。 + entry.getKey().processAll(entry.getValue()); + processed += entry.getValue().size(); } if (processed > 0) { // commitSync 放在批次末尾,保证同一个 poll 批次内的消息按 Kafka offset 一起确认。 diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqProperties.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqProperties.java index f4c46171..92aa9644 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqProperties.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqProperties.java @@ -93,6 +93,7 @@ public class SinkMqProperties { private int loopBackoffMillis = 1000; private String autoOffsetReset = "earliest"; private int maxPollRecords = 500; + private int maxPollIntervalMillis = 1800000; /** * Processor bean name -> Kafka binding。 * @@ -115,6 +116,8 @@ public class SinkMqProperties { public void setAutoOffsetReset(String autoOffsetReset) { this.autoOffsetReset = autoOffsetReset; } public int getMaxPollRecords() { return maxPollRecords; } public void setMaxPollRecords(int maxPollRecords) { this.maxPollRecords = maxPollRecords; } + public int getMaxPollIntervalMillis() { return maxPollIntervalMillis; } + public void setMaxPollIntervalMillis(int maxPollIntervalMillis) { this.maxPollIntervalMillis = maxPollIntervalMillis; } public Map getBindings() { return bindings; } public void setBindings(Map bindings) { this.bindings = bindings; } } diff --git a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerFactoryTest.java b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerFactoryTest.java index 291a1372..c1117ab2 100644 --- a/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerFactoryTest.java +++ b/modules/sinks/sink-mq/src/test/java/com/lingniu/ingest/sink/mq/KafkaEnvelopeConsumerFactoryTest.java @@ -40,6 +40,7 @@ class KafkaEnvelopeConsumerFactoryTest { .allSatisfy(p -> { assertThat(p.getProperty(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG)).isEqualTo("kafka-1:9092"); assertThat(p.get(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG)).isEqualTo(false); + assertThat(p.get(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG)).isEqualTo(1800000); }); } diff --git a/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineJdbcHistoryWriter.java b/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineJdbcHistoryWriter.java index 01acad61..008926af 100644 --- a/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineJdbcHistoryWriter.java +++ b/modules/sinks/tdengine-history-store/src/main/java/com/lingniu/ingest/tdenginehistory/TdengineJdbcHistoryWriter.java @@ -20,6 +20,7 @@ import org.slf4j.LoggerFactory; public final class TdengineJdbcHistoryWriter implements TdengineHistoryWriter { private static final Logger log = LoggerFactory.getLogger(TdengineJdbcHistoryWriter.class); + private static final int MAX_LITERAL_ROWS_PER_STATEMENT = 200; private final DataSource dataSource; private final TdengineHistorySchema schema; @@ -179,26 +180,39 @@ public final class TdengineJdbcHistoryWriter implements TdengineHistoryWriter { private static void executeLiteralBatch(Connection connection, LiteralBatch batch) throws SQLException { try (Statement statement = connection.createStatement()) { - for (List row : batch.rows()) { - statement.execute(toLiteralInsertSql(batch.insertSql(), row)); + List> rows = batch.rows(); + for (int start = 0; start < rows.size(); start += MAX_LITERAL_ROWS_PER_STATEMENT) { + int end = Math.min(start + MAX_LITERAL_ROWS_PER_STATEMENT, rows.size()); + statement.execute(toLiteralInsertSql(batch.insertSql(), rows.subList(start, end))); } } } - private static String toLiteralInsertSql(String insertSql, List values) { + private static String toLiteralInsertSql(String insertSql, List> rows) { int valuesIndex = insertSql.lastIndexOf("VALUES"); if (valuesIndex < 0) { throw new IllegalArgumentException("insert SQL must contain VALUES: " + insertSql); } - StringBuilder sql = new StringBuilder(insertSql.substring(0, valuesIndex)) - .append("VALUES ("); - for (int i = 0; i < values.size(); i++) { - if (i > 0) { - sql.append(", "); - } - sql.append(literal(values.get(i))); + if (rows == null || rows.isEmpty()) { + throw new IllegalArgumentException("literal rows must not be empty"); } - return sql.append(')').toString(); + StringBuilder sql = new StringBuilder(insertSql.substring(0, valuesIndex)) + .append("VALUES "); + for (int rowIndex = 0; rowIndex < rows.size(); rowIndex++) { + if (rowIndex > 0) { + sql.append(' '); + } + List values = rows.get(rowIndex); + sql.append('('); + for (int i = 0; i < values.size(); i++) { + if (i > 0) { + sql.append(", "); + } + sql.append(literal(values.get(i))); + } + sql.append(')'); + } + return sql.toString(); } private static String literal(Object value) { diff --git a/modules/sinks/tdengine-history-store/src/test/java/com/lingniu/ingest/tdenginehistory/TdengineJdbcHistoryWriterTest.java b/modules/sinks/tdengine-history-store/src/test/java/com/lingniu/ingest/tdenginehistory/TdengineJdbcHistoryWriterTest.java index da943639..f9e1823b 100644 --- a/modules/sinks/tdengine-history-store/src/test/java/com/lingniu/ingest/tdenginehistory/TdengineJdbcHistoryWriterTest.java +++ b/modules/sinks/tdengine-history-store/src/test/java/com/lingniu/ingest/tdenginehistory/TdengineJdbcHistoryWriterTest.java @@ -156,6 +156,25 @@ class TdengineJdbcHistoryWriterTest { assertThat(jdbc.commits).isEqualTo(1); } + @Test + void groupsLiteralRowsForSameChildTableIntoOneInsert() throws Exception { + RecordingJdbc jdbc = new RecordingJdbc(); + TdengineJdbcHistoryWriter writer = new TdengineJdbcHistoryWriter(jdbc.dataSource(), schema); + + writer.appendRawFrames(List.of( + rawFrame("frame-1", Instant.parse("2026-06-29T05:00:01Z")), + rawFrame("frame-2", Instant.parse("2026-06-29T05:00:02Z")))); + + List inserts = jdbc.executedSql.stream() + .filter(sql -> sql.startsWith("INSERT INTO raw_jt808_")) + .toList(); + assertThat(inserts).hasSize(1); + assertThat(inserts.getFirst()) + .contains("frame-1") + .contains("frame-2") + .contains(") ("); + } + private static TdengineRawFrameRow rawFrame(String frameId, Instant ts) { return rawFrame(frameId, ts, ""); }