refactor: expose kafka sink configuration

This commit is contained in:
lingniu
2026-07-01 10:50:50 +08:00
parent 303b2eb47e
commit 849e8f8a8b
67 changed files with 350 additions and 385 deletions

View File

@@ -1,4 +1,4 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import com.lingniu.ingest.api.event.AlarmPayload;
import com.lingniu.ingest.api.event.LocationPayload;
@@ -10,10 +10,10 @@ import com.lingniu.ingest.api.event.VehicleEvent;
import com.lingniu.ingest.api.event.VehicleEventTelemetrySnapshotMapper;
import com.lingniu.ingest.facts.FactIds;
import com.lingniu.ingest.facts.VehicleKey;
import com.lingniu.ingest.sink.mq.proto.ParseStatusProto;
import com.lingniu.ingest.sink.mq.proto.TelemetryField;
import com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot.Builder;
import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope;
import com.lingniu.ingest.sink.kafka.proto.ParseStatusProto;
import com.lingniu.ingest.sink.kafka.proto.TelemetryField;
import com.lingniu.ingest.sink.kafka.proto.TelemetrySnapshot.Builder;
import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope;
import java.security.MessageDigest;
import java.security.NoSuchAlgorithmException;
@@ -53,23 +53,23 @@ public final class EnvelopeMapper {
case VehicleEvent.Location l -> b.setLocation(buildLocation(l.payload()));
case VehicleEvent.Alarm a -> b.setAlarm(buildAlarm(a.payload()));
case VehicleEvent.Login lg -> b.setLogin(
com.lingniu.ingest.sink.mq.proto.LoginPayload.newBuilder()
com.lingniu.ingest.sink.kafka.proto.LoginPayload.newBuilder()
.setIccid(nullToEmpty(lg.iccid()))
.setProtocolVersion(nullToEmpty(lg.protocolVersion()))
.build());
case VehicleEvent.Logout ignored -> b.setLogout(
com.lingniu.ingest.sink.mq.proto.LogoutPayload.getDefaultInstance());
com.lingniu.ingest.sink.kafka.proto.LogoutPayload.getDefaultInstance());
case VehicleEvent.Heartbeat ignored -> b.setHeartbeat(
com.lingniu.ingest.sink.mq.proto.HeartbeatPayload.getDefaultInstance());
com.lingniu.ingest.sink.kafka.proto.HeartbeatPayload.getDefaultInstance());
case VehicleEvent.MediaMeta m -> b.setMediaMeta(
com.lingniu.ingest.sink.mq.proto.MediaMetaPayload.newBuilder()
com.lingniu.ingest.sink.kafka.proto.MediaMetaPayload.newBuilder()
.setMediaId(nullToEmpty(m.mediaId()))
.setMediaType(nullToEmpty(m.mediaType()))
.setSizeBytes(m.sizeBytes())
.setArchiveRef(nullToEmpty(m.archiveRef()))
.build());
case VehicleEvent.Passthrough p -> b.setPassthrough(
com.lingniu.ingest.sink.mq.proto.PassthroughPayload.newBuilder()
com.lingniu.ingest.sink.kafka.proto.PassthroughPayload.newBuilder()
.setPassthroughType(p.passthroughType())
.setData(com.google.protobuf.ByteString.copyFrom(
p.data() == null ? new byte[0] : p.data()))
@@ -88,13 +88,13 @@ public final class EnvelopeMapper {
b.putMetadata(RawArchiveKeys.META_KEY, key);
b.putMetadata(RawArchiveKeys.META_URI, uri);
b.putMetadata(RawArchiveKeys.META_EVENT_ID, ra.eventId());
b.setRawArchive(com.lingniu.ingest.sink.mq.proto.RawArchiveRef.newBuilder()
b.setRawArchive(com.lingniu.ingest.sink.kafka.proto.RawArchiveRef.newBuilder()
.setUri(uri)
.setChecksum(checksum)
.setSizeBytes(size)
.setParsedJson(nullToEmpty(ra.parsedJson()))
.build());
b.setRawFrameFact(com.lingniu.ingest.sink.mq.proto.RawFrameFactPayload.newBuilder()
b.setRawFrameFact(com.lingniu.ingest.sink.kafka.proto.RawFrameFactPayload.newBuilder()
.setFrameId(frameId)
.setVehicleKey(vehicleKey)
.setVin(nullToEmpty(ra.vin()))
@@ -114,8 +114,8 @@ public final class EnvelopeMapper {
return b.build();
}
private static com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot buildTelemetrySnapshot(TelemetrySnapshot snapshot) {
Builder b = com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot.newBuilder()
private static com.lingniu.ingest.sink.kafka.proto.TelemetrySnapshot buildTelemetrySnapshot(TelemetrySnapshot snapshot) {
Builder b = com.lingniu.ingest.sink.kafka.proto.TelemetrySnapshot.newBuilder()
.setEventType(snapshot.eventType())
.setRawArchiveUri(snapshot.rawArchiveUri());
for (TelemetryFieldValue field : snapshot.fields()) {
@@ -131,8 +131,8 @@ public final class EnvelopeMapper {
return b.build();
}
private static com.lingniu.ingest.sink.mq.proto.RealtimePayload buildRealtime(RealtimePayload p) {
var b = com.lingniu.ingest.sink.mq.proto.RealtimePayload.newBuilder();
private static com.lingniu.ingest.sink.kafka.proto.RealtimePayload buildRealtime(RealtimePayload p) {
var b = com.lingniu.ingest.sink.kafka.proto.RealtimePayload.newBuilder();
if (p.speedKmh() != null) b.setSpeedKmh(p.speedKmh());
if (p.totalMileageKm() != null) b.setTotalMileageKm(p.totalMileageKm());
if (p.batterySoc() != null) b.setBatterySoc(p.batterySoc());
@@ -159,8 +159,8 @@ public final class EnvelopeMapper {
return b.build();
}
private static com.lingniu.ingest.sink.mq.proto.LocationPayload buildLocation(LocationPayload p) {
return com.lingniu.ingest.sink.mq.proto.LocationPayload.newBuilder()
private static com.lingniu.ingest.sink.kafka.proto.LocationPayload buildLocation(LocationPayload p) {
return com.lingniu.ingest.sink.kafka.proto.LocationPayload.newBuilder()
.setLongitude(p.longitude())
.setLatitude(p.latitude())
.setAltitudeM(p.altitudeM())
@@ -171,8 +171,8 @@ public final class EnvelopeMapper {
.build();
}
private static com.lingniu.ingest.sink.mq.proto.AlarmPayload buildAlarm(AlarmPayload p) {
var b = com.lingniu.ingest.sink.mq.proto.AlarmPayload.newBuilder()
private static com.lingniu.ingest.sink.kafka.proto.AlarmPayload buildAlarm(AlarmPayload p) {
var b = com.lingniu.ingest.sink.kafka.proto.AlarmPayload.newBuilder()
.setLevel(p.level().name())
.setAlarmTypeCode(p.alarmTypeCode())
.setAlarmTypeName(nullToEmpty(p.alarmTypeName()));

View File

@@ -1,4 +1,4 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor;
import org.apache.kafka.clients.consumer.ConsumerConfig;
@@ -28,15 +28,15 @@ public final class KafkaEnvelopeConsumerFactory {
}
public List<KafkaEnvelopeConsumerWorker> createWorkers(Map<String, EnvelopeConsumerProcessor> processors,
SinkMqProperties props) {
Map<String, SinkMqProperties.Binding> bindings = effectiveBindings(props);
KafkaSinkProperties props) {
Map<String, KafkaSinkProperties.Binding> bindings = effectiveBindings(props);
List<KafkaEnvelopeConsumerWorker> workers = new ArrayList<>();
for (Map.Entry<String, SinkMqProperties.Binding> entry : bindings.entrySet()) {
for (Map.Entry<String, KafkaSinkProperties.Binding> entry : bindings.entrySet()) {
String processorBeanName = entry.getKey();
// binding key 必须和 Spring Bean 名一致这样配置只声明 topic/group
// 实际处理逻辑仍由各业务模块自己的 EnvelopeConsumerProcessor 承接
EnvelopeConsumerProcessor processor = processors.get(processorBeanName);
SinkMqProperties.Binding binding = entry.getValue();
KafkaSinkProperties.Binding binding = entry.getValue();
if (processor == null || binding == null || !binding.isEnabled()) {
continue;
}
@@ -57,15 +57,15 @@ public final class KafkaEnvelopeConsumerFactory {
return workers;
}
private Map<String, SinkMqProperties.Binding> effectiveBindings(SinkMqProperties props) {
Map<String, SinkMqProperties.Binding> configured = props.getConsumer().getBindings();
private Map<String, KafkaSinkProperties.Binding> effectiveBindings(KafkaSinkProperties props) {
Map<String, KafkaSinkProperties.Binding> configured = props.getConsumer().getBindings();
if (configured != null && !configured.isEmpty()) {
return configured;
}
// 默认绑定仅给未显式配置 bindings 的轻量运行时兜底
// 生产 history app 会显式绑定各协议 event/raw topic并写入 TDengine raw_frames/locations
SinkMqProperties.Topics topics = props.getTopics();
Map<String, SinkMqProperties.Binding> defaults = new LinkedHashMap<>();
KafkaSinkProperties.Topics topics = props.getTopics();
Map<String, KafkaSinkProperties.Binding> defaults = new LinkedHashMap<>();
defaults.put("eventHistoryEnvelopeConsumerProcessor", binding(
"vehicle-event-history",
topics.getRealtime(), topics.getLocation(), topics.getAlarm(), topics.getSession(), topics.getMediaMeta()));
@@ -78,8 +78,8 @@ public final class KafkaEnvelopeConsumerFactory {
return defaults;
}
private SinkMqProperties.Binding binding(String groupId, String... topics) {
SinkMqProperties.Binding binding = new SinkMqProperties.Binding();
private KafkaSinkProperties.Binding binding(String groupId, String... topics) {
KafkaSinkProperties.Binding binding = new KafkaSinkProperties.Binding();
binding.setGroupId(groupId);
binding.setTopics(List.of(topics));
return binding;
@@ -106,12 +106,12 @@ public final class KafkaEnvelopeConsumerFactory {
return List.copyOf(clean);
}
private Properties consumerProperties(SinkMqProperties props,
SinkMqProperties.Binding binding,
private Properties consumerProperties(KafkaSinkProperties props,
KafkaSinkProperties.Binding binding,
String processorBeanName,
int workerIndex,
int concurrency) {
SinkMqProperties.Consumer consumer = props.getConsumer();
KafkaSinkProperties.Consumer consumer = props.getConsumer();
Properties p = new Properties();
p.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, props.getBootstrapServers());
p.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
@@ -133,7 +133,7 @@ public final class KafkaEnvelopeConsumerFactory {
return concurrency <= 1 ? base : base + "-" + workerIndex;
}
private String groupId(SinkMqProperties.Binding binding, String processorBeanName) {
private String groupId(KafkaSinkProperties.Binding binding, String processorBeanName) {
if (binding.getGroupId() != null && !binding.getGroupId().isBlank()) {
return binding.getGroupId();
}

View File

@@ -1,4 +1,4 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

View File

@@ -1,4 +1,4 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor;
import com.lingniu.ingest.api.consumer.EnvelopeConsumerRecord;

View File

@@ -1,4 +1,4 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import com.lingniu.ingest.api.consumer.EnvelopeDeadLetterRecord;
import com.lingniu.ingest.api.consumer.EnvelopeDeadLetterSink;

View File

@@ -1,4 +1,4 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import com.lingniu.ingest.api.event.VehicleEvent;
import com.lingniu.ingest.api.sink.EventSink;

View File

@@ -1,4 +1,4 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import io.github.resilience4j.circuitbreaker.CircuitBreaker;
import io.github.resilience4j.circuitbreaker.CircuitBreakerConfig;
@@ -16,40 +16,31 @@ import java.time.Duration;
import java.util.Properties;
/**
* MQ Sink 自动装配
* Kafka Sink 自动装配
*
* <p>装配前提两个条件都要满足
* <ol>
* <li>{@code lingniu.ingest.sink.mq.enabled=true}默认 true <b>总开关</b>
* 设为 false 时本模块完全不装配ingest-core DisruptorEventBus 仍然运行但
* 没有 Kafka sink
* <li>{@code lingniu.ingest.sink.mq.type=kafka}默认 kafka 唯一生产 MQ 后端
* </ol>
*
* <p>Producer Consumer 是两个独立开关{@code sink.mq.enabled=true} 只表示可以创建
* Kafka producer/sink是否从 Kafka envelope 还要单独开启
* {@code lingniu.ingest.sink.mq.consumer.enabled=true} 并配置 bindings
* <p>{@code lingniu.ingest.sink.kafka.enabled=false} 时本模块完全不装配Producer Consumer
* 是两个独立开关消费端还需要开启 {@code lingniu.ingest.sink.kafka.consumer.enabled=true}
* 并配置 bindings
*/
@AutoConfiguration
@EnableConfigurationProperties(SinkMqProperties.class)
@ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq", name = "enabled", havingValue = "true", matchIfMissing = true)
public class SinkMqAutoConfiguration {
@EnableConfigurationProperties(KafkaSinkProperties.class)
@ConditionalOnProperty(prefix = "lingniu.ingest.sink.kafka", name = "enabled", havingValue = "true", matchIfMissing = true)
public class KafkaSinkAutoConfiguration {
@Bean
@ConditionalOnMissingBean
public EnvelopeMapper envelopeMapper(SinkMqProperties props) {
public EnvelopeMapper envelopeMapper(KafkaSinkProperties props) {
return new EnvelopeMapper(props.getNodeId());
}
@Bean
@ConditionalOnMissingBean
public TopicRouter topicRouter(SinkMqProperties props) {
public TopicRouter topicRouter(KafkaSinkProperties props) {
return new TopicRouter(props.getTopics());
}
@Bean
@ConditionalOnMissingBean
@ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq", name = "type", havingValue = "kafka", matchIfMissing = true)
public CircuitBreaker kafkaSinkCircuitBreaker() {
return CircuitBreaker.of("kafka-sink", CircuitBreakerConfig.custom()
.slidingWindowSize(100)
@@ -61,8 +52,7 @@ public class SinkMqAutoConfiguration {
@Bean(destroyMethod = "close")
@ConditionalOnMissingBean
@ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq", name = "type", havingValue = "kafka", matchIfMissing = true)
public KafkaProducer<String, byte[]> kafkaProducer(SinkMqProperties props) {
public KafkaProducer<String, byte[]> kafkaProducer(KafkaSinkProperties props) {
Properties p = new Properties();
p.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, props.getBootstrapServers());
p.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
@@ -77,11 +67,10 @@ public class SinkMqAutoConfiguration {
@Bean(destroyMethod = "close")
@ConditionalOnMissingBean
@ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq", name = "type", havingValue = "kafka", matchIfMissing = true)
public KafkaEventSink kafkaEventSink(KafkaProducer<String, byte[]> producer,
EnvelopeMapper mapper,
TopicRouter router,
SinkMqProperties props,
KafkaSinkProperties props,
CircuitBreaker breaker) {
// KafkaEventSink EventBus 的生产端出口它不会启动任何 Kafka 消费线程
return new KafkaEventSink(producer, mapper, router, props.getTopics().getDlq(), breaker);
@@ -89,9 +78,8 @@ public class SinkMqAutoConfiguration {
@Bean
@ConditionalOnMissingBean
@ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq", name = "type", havingValue = "kafka", matchIfMissing = true)
public KafkaEnvelopeDeadLetterSink kafkaEnvelopeDeadLetterSink(KafkaProducer<String, byte[]> producer,
SinkMqProperties props) {
KafkaSinkProperties props) {
return new KafkaEnvelopeDeadLetterSink(producer, props.getTopics().getDlq());
}

View File

@@ -1,4 +1,4 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor;
import org.springframework.beans.factory.ListableBeanFactory;
@@ -13,15 +13,15 @@ import java.time.Duration;
import java.util.List;
import java.util.Map;
@AutoConfiguration(after = SinkMqAutoConfiguration.class)
@AutoConfiguration(after = KafkaSinkAutoConfiguration.class)
@AutoConfigureAfter(name = {
"com.lingniu.ingest.eventhistory.config.EventHistoryAutoConfiguration",
"com.lingniu.ingest.vehiclestate.config.VehicleStateAutoConfiguration",
"com.lingniu.ingest.vehiclestat.config.VehicleStatAutoConfiguration"
})
@EnableConfigurationProperties(SinkMqProperties.class)
@ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq", name = "enabled", havingValue = "true", matchIfMissing = true)
public class SinkMqConsumerAutoConfiguration {
@EnableConfigurationProperties(KafkaSinkProperties.class)
@ConditionalOnProperty(prefix = "lingniu.ingest.sink.kafka", name = "enabled", havingValue = "true", matchIfMissing = true)
public class KafkaSinkConsumerAutoConfiguration {
@Bean
@ConditionalOnMissingBean
@@ -31,10 +31,10 @@ public class SinkMqConsumerAutoConfiguration {
@Bean
@ConditionalOnMissingBean
@ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq.consumer", name = "enabled", havingValue = "true")
@ConditionalOnProperty(prefix = "lingniu.ingest.sink.kafka.consumer", name = "enabled", havingValue = "true")
public KafkaEnvelopeConsumerRunner kafkaEnvelopeConsumerRunner(ListableBeanFactory beanFactory,
KafkaEnvelopeConsumerFactory consumerFactory,
SinkMqProperties props) {
KafkaSinkProperties props) {
return new KafkaEnvelopeConsumerRunner(
() -> createWorkers(beanFactory, consumerFactory, props),
Duration.ofMillis(props.getConsumer().getPollTimeoutMillis()),
@@ -44,7 +44,7 @@ public class SinkMqConsumerAutoConfiguration {
private List<KafkaEnvelopeConsumerWorker> createWorkers(ListableBeanFactory beanFactory,
KafkaEnvelopeConsumerFactory consumerFactory,
SinkMqProperties props) {
KafkaSinkProperties props) {
Map<String, EnvelopeConsumerProcessor> processors = beanFactory.getBeansOfType(EnvelopeConsumerProcessor.class);
return consumerFactory.createWorkers(processors, props);
}

View File

@@ -1,26 +1,23 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import org.springframework.boot.context.properties.ConfigurationProperties;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
@ConfigurationProperties(prefix = "lingniu.ingest.sink.mq")
public class SinkMqProperties {
@ConfigurationProperties(prefix = "lingniu.ingest.sink.kafka")
public class KafkaSinkProperties {
/**
* MQ Sink 总开关默认 {@code true}
* 设为 {@code false} {@link SinkMqAutoConfiguration} 完全不装配任何 BeanKafka Producer
* Kafka Sink 总开关默认 {@code true}
* 设为 {@code false} {@link KafkaSinkAutoConfiguration} 完全不装配任何 BeanKafka Producer
* EnvelopeMapperTopicRouterKafkaEventSink 都不会创建ingest-core DisruptorEventBus
* 仍然正常运行但没有外部 sink事件落地到 sink-archive Noop 吞掉
*/
private boolean enabled = true;
/** MQ 后端类型。生产链路只允许 kafka。 */
private String type = "kafka";
/** Kafka bootstrap servers生产环境应通过环境变量覆盖不建议使用默认开发地址。 */
private String bootstrapServers = "114.55.58.251:9092";
private String compressionType = "zstd";
@@ -35,14 +32,6 @@ public class SinkMqProperties {
public boolean isEnabled() { return enabled; }
public void setEnabled(boolean enabled) { this.enabled = enabled; }
public String getType() { return type; }
public void setType(String type) {
String value = type == null || type.isBlank() ? "kafka" : type.trim().toLowerCase(Locale.ROOT);
if (!"kafka".equals(value)) {
throw new IllegalStateException("sink.mq.type supports only kafka; configured value: " + type);
}
this.type = value;
}
public String getBootstrapServers() { return bootstrapServers; }
public void setBootstrapServers(String bootstrapServers) { this.bootstrapServers = bootstrapServers; }
public String getCompressionType() { return compressionType; }

View File

@@ -1,4 +1,4 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import com.lingniu.ingest.api.event.VehicleEvent;
@@ -7,9 +7,9 @@ import com.lingniu.ingest.api.event.VehicleEvent;
*/
public final class TopicRouter {
private final SinkMqProperties.Topics topics;
private final KafkaSinkProperties.Topics topics;
public TopicRouter(SinkMqProperties.Topics topics) {
public TopicRouter(KafkaSinkProperties.Topics topics) {
this.topics = topics;
}

View File

@@ -1,9 +1,9 @@
syntax = "proto3";
package com.lingniu.ingest.sink.mq.proto;
package com.lingniu.ingest.sink.kafka.proto;
option java_multiple_files = true;
option java_package = "com.lingniu.ingest.sink.mq.proto";
option java_package = "com.lingniu.ingest.sink.kafka.proto";
option java_outer_classname = "VehicleEnvelopeProto";
// 统一消息外壳:所有 Topic 共用此 Envelopepayload 通过 oneof 区分具体事件类型。

View File

@@ -1,2 +1,2 @@
com.lingniu.ingest.sink.mq.SinkMqAutoConfiguration
com.lingniu.ingest.sink.mq.SinkMqConsumerAutoConfiguration
com.lingniu.ingest.sink.kafka.KafkaSinkAutoConfiguration
com.lingniu.ingest.sink.kafka.KafkaSinkConsumerAutoConfiguration

View File

@@ -1,10 +1,10 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import com.lingniu.ingest.api.ProtocolId;
import com.lingniu.ingest.api.event.RawArchiveKeys;
import com.lingniu.ingest.api.event.RealtimePayload;
import com.lingniu.ingest.api.event.VehicleEvent;
import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope;
import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope;
import org.junit.jupiter.api.Test;
import java.time.Instant;

View File

@@ -1,4 +1,4 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor;
import com.lingniu.ingest.api.consumer.EnvelopeIngestResult;
@@ -23,7 +23,7 @@ class KafkaEnvelopeConsumerFactoryTest {
created.add(props);
return new MockConsumer<>(OffsetResetStrategy.EARLIEST);
});
SinkMqProperties props = new SinkMqProperties();
KafkaSinkProperties props = new KafkaSinkProperties();
props.setBootstrapServers("kafka-1:9092");
EnvelopeConsumerProcessor stateProcessor = processor();
EnvelopeConsumerProcessor statProcessor = processor();
@@ -51,7 +51,7 @@ class KafkaEnvelopeConsumerFactoryTest {
created.add(props);
return new MockConsumer<>(OffsetResetStrategy.EARLIEST);
});
SinkMqProperties props = new SinkMqProperties();
KafkaSinkProperties props = new KafkaSinkProperties();
props.getConsumer().setConcurrency(3);
EnvelopeConsumerProcessor stateProcessor = processor();

View File

@@ -1,4 +1,4 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import org.junit.jupiter.api.Test;

View File

@@ -1,4 +1,4 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import com.lingniu.ingest.api.consumer.EnvelopeConsumerRecord;
import com.lingniu.ingest.api.consumer.EnvelopeBatchIngestor;

View File

@@ -1,4 +1,4 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import com.lingniu.ingest.api.consumer.EnvelopeDeadLetterRecord;
import com.lingniu.ingest.api.consumer.EnvelopeIngestResult;

View File

@@ -1,4 +1,4 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import com.lingniu.ingest.api.ProtocolId;
import com.lingniu.ingest.api.event.LocationPayload;
@@ -26,7 +26,7 @@ class KafkaEventSinkTest {
KafkaEventSink sink = new KafkaEventSink(
null,
new EnvelopeMapper("node-1"),
new TopicRouter(new SinkMqProperties.Topics()),
new TopicRouter(new KafkaSinkProperties.Topics()),
"vehicle.dlq.gb32960.v1",
CircuitBreaker.ofDefaults("kafka-test"));
@@ -46,7 +46,7 @@ class KafkaEventSinkTest {
KafkaEventSink sink = new KafkaEventSink(
producer,
new EnvelopeMapper("node-1"),
new TopicRouter(new SinkMqProperties.Topics()),
new TopicRouter(new KafkaSinkProperties.Topics()),
"vehicle.dlq.gb32960.v1",
CircuitBreaker.ofDefaults("kafka-test"));

View File

@@ -1,4 +1,4 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor;
import com.lingniu.ingest.api.consumer.EnvelopeIngestResult;
@@ -15,26 +15,25 @@ import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.Mockito.mock;
@ExtendWith(OutputCaptureExtension.class)
class SinkMqConsumerAutoConfigurationTest {
class KafkaSinkConsumerAutoConfigurationTest {
private final ApplicationContextRunner contextRunner = new ApplicationContextRunner()
.withUserConfiguration(SinkMqAutoConfiguration.class, SinkMqConsumerAutoConfiguration.class)
.withUserConfiguration(KafkaSinkAutoConfiguration.class, KafkaSinkConsumerAutoConfiguration.class)
.withAllowBeanDefinitionOverriding(true)
.withBean("vehicleStateEnvelopeConsumerProcessor", EnvelopeConsumerProcessor.class,
() -> new EnvelopeConsumerProcessor(
"vehicle-state",
bytes -> EnvelopeIngestResult.processed("evt-1", "VIN001"),
record -> {}))
.withBean("kafkaProducer", KafkaProducer.class, SinkMqConsumerAutoConfigurationTest::kafkaProducer)
.withBean("kafkaProducer", KafkaProducer.class, KafkaSinkConsumerAutoConfigurationTest::kafkaProducer)
.withBean(KafkaEnvelopeConsumerFactory.class,
() -> new KafkaEnvelopeConsumerFactory(props -> new MockConsumer<>(OffsetResetStrategy.EARLIEST)))
.withPropertyValues(
"lingniu.ingest.sink.mq.enabled=true",
"lingniu.ingest.sink.mq.type=kafka",
"lingniu.ingest.sink.mq.consumer.enabled=true",
"lingniu.ingest.sink.mq.consumer.auto-startup=false",
"lingniu.ingest.sink.mq.consumer.bindings.vehicleStateEnvelopeConsumerProcessor.group-id=vehicle-state",
"lingniu.ingest.sink.mq.consumer.bindings.vehicleStateEnvelopeConsumerProcessor.topics[0]=vehicle.realtime");
"lingniu.ingest.sink.kafka.enabled=true",
"lingniu.ingest.sink.kafka.consumer.enabled=true",
"lingniu.ingest.sink.kafka.consumer.auto-startup=false",
"lingniu.ingest.sink.kafka.consumer.bindings.vehicleStateEnvelopeConsumerProcessor.group-id=vehicle-state",
"lingniu.ingest.sink.kafka.consumer.bindings.vehicleStateEnvelopeConsumerProcessor.topics[0]=vehicle.realtime");
@Test
void createsKafkaEnvelopeConsumerRunnerWhenConsumerBindingIsConfigured(CapturedOutput output) {

View File

@@ -0,0 +1,26 @@
package com.lingniu.ingest.sink.kafka;
import org.junit.jupiter.api.Test;
import java.util.Arrays;
import static org.assertj.core.api.Assertions.assertThat;
class KafkaSinkPropertiesTest {
@Test
void defaultKafkaBrokerUsesProductionAddressButCanStillBeOverriddenByConfigBinding() {
KafkaSinkProperties props = new KafkaSinkProperties();
assertThat(props.getBootstrapServers()).isEqualTo("114.55.58.251:9092");
props.setBootstrapServers("kafka.internal:9092");
assertThat(props.getBootstrapServers()).isEqualTo("kafka.internal:9092");
}
@Test
void exposesNoBackendTypeBecauseKafkaIsTheOnlySink() {
assertThat(Arrays.stream(KafkaSinkProperties.class.getMethods()).map(method -> method.getName()))
.doesNotContain("getType", "setType");
}
}

View File

@@ -1,4 +1,4 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import com.lingniu.ingest.api.ProtocolId;
import com.lingniu.ingest.api.event.AlarmPayload;
@@ -18,14 +18,14 @@ class TopicRouterTest {
@Test
void rawArchiveRoutesToVersionedGb32960RawTopicByDefault() {
TopicRouter router = new TopicRouter(new SinkMqProperties.Topics());
TopicRouter router = new TopicRouter(new KafkaSinkProperties.Topics());
assertThat(router.route(rawArchive())).isEqualTo("vehicle.raw.gb32960.v1");
}
@Test
void normalizedGb32960EventsRouteToVersionedEventTopicByDefault() {
TopicRouter router = new TopicRouter(new SinkMqProperties.Topics());
TopicRouter router = new TopicRouter(new KafkaSinkProperties.Topics());
assertThat(router.route(realtime())).isEqualTo("vehicle.event.gb32960.v1");
assertThat(router.route(location())).isEqualTo("vehicle.event.gb32960.v1");
@@ -37,7 +37,7 @@ class TopicRouterTest {
@Test
void dlqDefaultsToVersionedGb32960DlqTopic() {
SinkMqProperties.Topics topics = new SinkMqProperties.Topics();
KafkaSinkProperties.Topics topics = new KafkaSinkProperties.Topics();
assertThat(topics.getDlq()).isEqualTo("vehicle.dlq.gb32960.v1");
}

View File

@@ -1,9 +1,9 @@
package com.lingniu.ingest.sink.mq;
package com.lingniu.ingest.sink.kafka;
import com.lingniu.ingest.sink.mq.proto.DecodedFactPayload;
import com.lingniu.ingest.sink.mq.proto.ParseStatusProto;
import com.lingniu.ingest.sink.mq.proto.RawFrameFactPayload;
import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope;
import com.lingniu.ingest.sink.kafka.proto.DecodedFactPayload;
import com.lingniu.ingest.sink.kafka.proto.ParseStatusProto;
import com.lingniu.ingest.sink.kafka.proto.RawFrameFactPayload;
import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope;
import org.junit.jupiter.api.Test;
import static org.assertj.core.api.Assertions.assertThat;

View File

@@ -1,29 +0,0 @@
package com.lingniu.ingest.sink.mq;
import org.junit.jupiter.api.Test;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
class SinkMqPropertiesTest {
@Test
void defaultKafkaBrokerUsesProductionAddressButCanStillBeOverriddenByConfigBinding() {
SinkMqProperties props = new SinkMqProperties();
assertThat(props.getBootstrapServers()).isEqualTo("114.55.58.251:9092");
props.setBootstrapServers("kafka.internal:9092");
assertThat(props.getBootstrapServers()).isEqualTo("kafka.internal:9092");
}
@Test
void rejectsUnsupportedMqTypeInsteadOfSilentlyDisablingKafkaBeans() {
SinkMqProperties props = new SinkMqProperties();
assertThatThrownBy(() -> props.setType("rocketmq"))
.isInstanceOf(IllegalStateException.class)
.hasMessageContaining("sink.mq.type supports only kafka")
.hasMessageContaining("rocketmq");
}
}

View File

@@ -2,10 +2,10 @@ package com.lingniu.ingest.tdenginehistory;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.lingniu.ingest.sink.mq.proto.LocationPayload;
import com.lingniu.ingest.sink.mq.proto.RawFrameFactPayload;
import com.lingniu.ingest.sink.mq.proto.TelemetryField;
import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope;
import com.lingniu.ingest.sink.kafka.proto.LocationPayload;
import com.lingniu.ingest.sink.kafka.proto.RawFrameFactPayload;
import com.lingniu.ingest.sink.kafka.proto.TelemetryField;
import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope;
import java.time.Instant;
import java.util.ArrayList;

View File

@@ -1,12 +1,12 @@
package com.lingniu.ingest.tdenginehistory;
import com.lingniu.ingest.sink.mq.proto.LocationPayload;
import com.lingniu.ingest.sink.mq.proto.ParseStatusProto;
import com.lingniu.ingest.sink.mq.proto.RawArchiveRef;
import com.lingniu.ingest.sink.mq.proto.RawFrameFactPayload;
import com.lingniu.ingest.sink.mq.proto.TelemetryField;
import com.lingniu.ingest.sink.mq.proto.TelemetrySnapshot;
import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope;
import com.lingniu.ingest.sink.kafka.proto.LocationPayload;
import com.lingniu.ingest.sink.kafka.proto.ParseStatusProto;
import com.lingniu.ingest.sink.kafka.proto.RawArchiveRef;
import com.lingniu.ingest.sink.kafka.proto.RawFrameFactPayload;
import com.lingniu.ingest.sink.kafka.proto.TelemetryField;
import com.lingniu.ingest.sink.kafka.proto.TelemetrySnapshot;
import com.lingniu.ingest.sink.kafka.proto.VehicleEnvelope;
import org.junit.jupiter.api.Test;
import java.time.Instant;