fix: keep jt808 raw history consumer current
This commit is contained in:
@@ -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<KafkaEnvelopeConsumerWorker> 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");
|
||||
|
||||
@@ -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}
|
||||
|
||||
@@ -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}");
|
||||
}
|
||||
|
||||
|
||||
@@ -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<EnvelopeIngestResult> processAll(List<EnvelopeConsumerRecord> records) {
|
||||
if (records == null) {
|
||||
throw new IllegalArgumentException("records must not be null");
|
||||
}
|
||||
if (records.isEmpty()) {
|
||||
return List.of();
|
||||
}
|
||||
List<byte[]> 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<EnvelopeIngestResult> 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) {
|
||||
|
||||
@@ -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<EnvelopeIngestResult> tryIngestAll(List<byte[]> kafkaValues) {
|
||||
if (kafkaValues == null || kafkaValues.isEmpty()) {
|
||||
return List.of();
|
||||
}
|
||||
List<EnvelopeIngestResult> results = new ArrayList<>(kafkaValues.size());
|
||||
for (byte[] kafkaValue : kafkaValues) {
|
||||
results.add(tryIngest(kafkaValue));
|
||||
}
|
||||
return results;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,31 +48,58 @@ 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) {
|
||||
return tryIngestAll(List.of(kafkaValue)).getFirst();
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<EnvelopeIngestResult> tryIngestAll(List<byte[]> kafkaValues) {
|
||||
if (kafkaValues == null || kafkaValues.isEmpty()) {
|
||||
return List.of();
|
||||
}
|
||||
List<EnvelopeIngestResult> results = new ArrayList<>(kafkaValues.size());
|
||||
List<BatchEntry> valid = new ArrayList<>(kafkaValues.size());
|
||||
for (byte[] kafkaValue : kafkaValues) {
|
||||
VehicleEnvelope envelope = null;
|
||||
try {
|
||||
envelope = parse(kafkaValue);
|
||||
EventFileRecord record = mapper.toRecord(envelope);
|
||||
store.append(record);
|
||||
writeTdengineFacts(envelope);
|
||||
return EnvelopeIngestResult.stored(record.eventId(), record.vin());
|
||||
results.add(null);
|
||||
valid.add(new BatchEntry(results.size() - 1, envelope, record));
|
||||
} catch (IllegalArgumentException ex) {
|
||||
// 协议不可解析或缺少必要字段属于不可重试问题,避免 Kafka consumer 无限重放。
|
||||
return envelope == null
|
||||
results.add(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());
|
||||
: 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) {
|
||||
if (kafkaValue == null || kafkaValue.length == 0) {
|
||||
@@ -84,19 +112,25 @@ public final class EventHistoryEnvelopeIngestor implements EnvelopeIngestor {
|
||||
}
|
||||
}
|
||||
|
||||
private void writeTdengineFacts(VehicleEnvelope envelope) throws IOException {
|
||||
private void writeTdengineFacts(List<VehicleEnvelope> 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);
|
||||
}
|
||||
|
||||
@@ -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<EnvelopeIngestResult> 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<EventFileRecord> records = new ArrayList<>();
|
||||
private int appendAllCalls;
|
||||
|
||||
@Override
|
||||
public void appendAll(List<EventFileRecord> records) {
|
||||
appendAllCalls++;
|
||||
this.records.addAll(records);
|
||||
}
|
||||
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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<String, byte[]> records = consumer.poll(timeout == null ? Duration.ZERO : timeout);
|
||||
int processed = 0;
|
||||
Map<EnvelopeConsumerProcessor, List<EnvelopeConsumerRecord>> batches = new LinkedHashMap<>();
|
||||
for (ConsumerRecord<String, byte[]> 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<EnvelopeConsumerProcessor, List<EnvelopeConsumerRecord>> entry : batches.entrySet()) {
|
||||
// EnvelopeConsumerProcessor 内部会把解析或业务错误转成 DLQ 记录,
|
||||
// 这里保持 Kafka worker 的职责单一:轮询、批量分发、成功后提交 offset。
|
||||
entry.getKey().processAll(entry.getValue());
|
||||
processed += entry.getValue().size();
|
||||
}
|
||||
if (processed > 0) {
|
||||
// commitSync 放在批次末尾,保证同一个 poll 批次内的消息按 Kafka offset 一起确认。
|
||||
|
||||
@@ -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<String, Binding> getBindings() { return bindings; }
|
||||
public void setBindings(Map<String, Binding> bindings) { this.bindings = bindings; }
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -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<Object> row : batch.rows()) {
|
||||
statement.execute(toLiteralInsertSql(batch.insertSql(), row));
|
||||
List<List<Object>> 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<Object> values) {
|
||||
private static String toLiteralInsertSql(String insertSql, List<List<Object>> rows) {
|
||||
int valuesIndex = insertSql.lastIndexOf("VALUES");
|
||||
if (valuesIndex < 0) {
|
||||
throw new IllegalArgumentException("insert SQL must contain VALUES: " + insertSql);
|
||||
}
|
||||
if (rows == null || rows.isEmpty()) {
|
||||
throw new IllegalArgumentException("literal rows must not be empty");
|
||||
}
|
||||
StringBuilder sql = new StringBuilder(insertSql.substring(0, valuesIndex))
|
||||
.append("VALUES (");
|
||||
.append("VALUES ");
|
||||
for (int rowIndex = 0; rowIndex < rows.size(); rowIndex++) {
|
||||
if (rowIndex > 0) {
|
||||
sql.append(' ');
|
||||
}
|
||||
List<Object> 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)));
|
||||
}
|
||||
return sql.append(')').toString();
|
||||
sql.append(')');
|
||||
}
|
||||
return sql.toString();
|
||||
}
|
||||
|
||||
private static String literal(Object value) {
|
||||
|
||||
@@ -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<String> 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, "");
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user