From 93f668deeb0fd8f7496dcd171040839dd7e360fb Mon Sep 17 00:00:00 2001 From: pthirun Date: Tue, 28 Jul 2026 17:19:23 -0700 Subject: [PATCH 1/3] [vpj] Fix NorthGuard repush ZSTD dictionary routing Preserve the configured pub-sub adapter and NorthGuard routing properties when fetching repush dictionaries, while overriding both broker settings to the source broker. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- .../linkedin/venice/hadoop/VenicePushJob.java | 35 ++++++----- .../hadoop/VenicePushJobRepushTest.java | 58 +++++++++++++++++++ 2 files changed, 80 insertions(+), 13 deletions(-) diff --git a/clients/venice-push-job/src/main/java/com/linkedin/venice/hadoop/VenicePushJob.java b/clients/venice-push-job/src/main/java/com/linkedin/venice/hadoop/VenicePushJob.java index 3f5968e6f94..7e12048b853 100755 --- a/clients/venice-push-job/src/main/java/com/linkedin/venice/hadoop/VenicePushJob.java +++ b/clients/venice-push-job/src/main/java/com/linkedin/venice/hadoop/VenicePushJob.java @@ -6,6 +6,7 @@ import static com.linkedin.venice.ConfigKeys.KAFKA_PRODUCER_REQUEST_TIMEOUT_MS; import static com.linkedin.venice.ConfigKeys.KAFKA_PRODUCER_RETRIES_CONFIG; import static com.linkedin.venice.ConfigKeys.MULTI_REGION; +import static com.linkedin.venice.ConfigKeys.PUBSUB_BROKER_ADDRESS; import static com.linkedin.venice.ConfigKeys.VENICE_PARTITIONERS; import static com.linkedin.venice.VeniceConstants.DEFAULT_SSL_FACTORY_CLASS_NAME; import static com.linkedin.venice.status.BatchJobHeartbeatConfigs.HEARTBEAT_ENABLED_CONFIG; @@ -866,13 +867,8 @@ public void run() { if (pushJobSetting.isSourceKafka) { if (pushJobSetting.sourceVersionCompressionStrategy == CompressionStrategy.ZSTD_WITH_DICT) { LOGGER.info("Source version uses ZSTD_WITH_DICT. Fetching source dictionary."); - Properties kafkaConsumerProperties = new Properties(); - if (pushJobSetting.enableSSL) { - kafkaConsumerProperties.putAll(this.sslProperties.get()); - } - kafkaConsumerProperties.setProperty(KAFKA_BOOTSTRAP_SERVERS, pushJobSetting.repushSourcePubsubBroker); ByteBuffer sourceDict = DictionaryUtils - .readDictionaryFromKafka(pushJobSetting.kafkaInputTopic, new VeniceProperties(kafkaConsumerProperties)); + .readDictionaryFromKafka(pushJobSetting.kafkaInputTopic, getSourceDictionaryConsumerProperties()); if (sourceDict != null) { pushJobSetting.sourceDictionary = ByteUtils.extractByteArray(sourceDict); } @@ -1666,6 +1662,25 @@ private Optional getCompressionDictionary() throws VeniceException { return Optional.of(emptyPushZstdDictionary.get()); } + private VeniceProperties getSourceDictionaryConsumerProperties() { + return buildSourceDictionaryConsumerProperties( + props, + pushJobSetting.enableSSL ? sslProperties.get() : new Properties(), + pushJobSetting.repushSourcePubsubBroker); + } + + @VisibleForTesting + static VeniceProperties buildSourceDictionaryConsumerProperties( + VeniceProperties jobProperties, + Properties sslProperties, + String sourcePubsubBroker) { + Properties consumerProperties = jobProperties.toProperties(); + consumerProperties.putAll(sslProperties); + consumerProperties.setProperty(PUBSUB_BROKER_ADDRESS, sourcePubsubBroker); + consumerProperties.setProperty(KAFKA_BOOTSTRAP_SERVERS, sourcePubsubBroker); + return new VeniceProperties(consumerProperties); + } + private ByteBuffer fetchOrBuildCompressionDictionary() throws VeniceException { // Prepare the param builder, which can be used by different scenarios. KafkaInputDictTrainer.ParamBuilder paramBuilder = new KafkaInputDictTrainer.ParamBuilder() @@ -1691,14 +1706,8 @@ private ByteBuffer fetchOrBuildCompressionDictionary() throws VeniceException { return ByteBuffer.wrap(dictTrainer.trainDict()); } else { LOGGER.info("Reading Zstd dictionary from input topic: {}", pushJobSetting.kafkaInputTopic); - // set up ssl properties and kafka consumer properties - Properties kafkaConsumerProperties = new Properties(); - if (pushJobSetting.enableSSL) { - kafkaConsumerProperties.putAll(this.sslProperties.get()); - } - kafkaConsumerProperties.setProperty(KAFKA_BOOTSTRAP_SERVERS, pushJobSetting.repushSourcePubsubBroker); return DictionaryUtils - .readDictionaryFromKafka(pushJobSetting.kafkaInputTopic, new VeniceProperties(kafkaConsumerProperties)); + .readDictionaryFromKafka(pushJobSetting.kafkaInputTopic, getSourceDictionaryConsumerProperties()); } } LOGGER.info( diff --git a/clients/venice-push-job/src/test/java/com/linkedin/venice/hadoop/VenicePushJobRepushTest.java b/clients/venice-push-job/src/test/java/com/linkedin/venice/hadoop/VenicePushJobRepushTest.java index f731eafa834..725d308c7cc 100644 --- a/clients/venice-push-job/src/test/java/com/linkedin/venice/hadoop/VenicePushJobRepushTest.java +++ b/clients/venice-push-job/src/test/java/com/linkedin/venice/hadoop/VenicePushJobRepushTest.java @@ -1,5 +1,10 @@ package com.linkedin.venice.hadoop; +import static com.linkedin.venice.ConfigKeys.KAFKA_BOOTSTRAP_SERVERS; +import static com.linkedin.venice.ConfigKeys.PUBSUB_BROKER_ADDRESS; +import static com.linkedin.venice.ConfigKeys.PUBSUB_CONSUMER_ADAPTER_FACTORY_CLASS; +import static com.linkedin.venice.ConfigKeys.PUBSUB_SECURITY_PROTOCOL; +import static com.linkedin.venice.ConfigKeys.PUBSUB_TYPE_ID_TO_POSITION_CLASS_NAME_MAP; import static com.linkedin.venice.vpj.VenicePushJobConstants.ALLOW_REGULAR_PUSH_WITH_TTL_REPUSH; import static com.linkedin.venice.vpj.VenicePushJobConstants.COMPLIANCE_PUSH; import static com.linkedin.venice.vpj.VenicePushJobConstants.KAFKA_INPUT_MAX_RECORDS_PER_MAPPER; @@ -26,6 +31,7 @@ import com.linkedin.venice.meta.StoreInfo; import com.linkedin.venice.meta.Version; import com.linkedin.venice.utils.Time; +import com.linkedin.venice.utils.VeniceProperties; import java.util.HashMap; import java.util.Map; import java.util.Properties; @@ -42,6 +48,58 @@ */ public class VenicePushJobRepushTest extends VenicePushJobTestBase { + @Test + public void testSourceDictionaryConsumerPropertiesRetainPubSubConfigAndOverrideBrokers() { + Properties jobProperties = new Properties(); + jobProperties.setProperty( + PUBSUB_CONSUMER_ADAPTER_FACTORY_CLASS, + "com.linkedin.venice.pubsub.adapter.xinfra.consumer.XcConsumerAdapterFactory"); + jobProperties.setProperty( + PUBSUB_TYPE_ID_TO_POSITION_CLASS_NAME_MAP, + "1:com.linkedin.venice.pubsub.adapter.xinfra.XinfraPosition"); + jobProperties.setProperty("xc.pubsub.broker.url.to.region.name.map", "northguard:ei4"); + jobProperties.setProperty(PUBSUB_BROKER_ADDRESS, "destination-broker"); + jobProperties.setProperty(KAFKA_BOOTSTRAP_SERVERS, "legacy-broker"); + + VeniceProperties consumerProperties = VenicePushJob.buildSourceDictionaryConsumerProperties( + new VeniceProperties(jobProperties), + new Properties(), + "source-broker"); + + assertEquals( + consumerProperties.getString(PUBSUB_CONSUMER_ADAPTER_FACTORY_CLASS), + "com.linkedin.venice.pubsub.adapter.xinfra.consumer.XcConsumerAdapterFactory"); + assertEquals( + consumerProperties.getString(PUBSUB_TYPE_ID_TO_POSITION_CLASS_NAME_MAP), + "1:com.linkedin.venice.pubsub.adapter.xinfra.XinfraPosition"); + assertEquals(consumerProperties.getString("xc.pubsub.broker.url.to.region.name.map"), "northguard:ei4"); + assertEquals(consumerProperties.getString(PUBSUB_BROKER_ADDRESS), "source-broker"); + assertEquals(consumerProperties.getString(KAFKA_BOOTSTRAP_SERVERS), "source-broker"); + assertEquals( + jobProperties.getProperty(PUBSUB_BROKER_ADDRESS), + "destination-broker", + "Building consumer properties must not mutate the job properties"); + } + + @Test + public void testSourceDictionaryConsumerPropertiesApplySslOverrides() { + Properties jobProperties = new Properties(); + jobProperties.setProperty(PUBSUB_SECURITY_PROTOCOL, "PLAINTEXT"); + jobProperties.setProperty("ssl.keystore.location", "stale-keystore"); + jobProperties.setProperty("xc.tls.key.store.type", "PKCS12"); + + Properties sslProperties = new Properties(); + sslProperties.setProperty(PUBSUB_SECURITY_PROTOCOL, "SSL"); + sslProperties.setProperty("ssl.keystore.location", "credential-keystore"); + + VeniceProperties consumerProperties = VenicePushJob + .buildSourceDictionaryConsumerProperties(new VeniceProperties(jobProperties), sslProperties, "source-broker"); + + assertEquals(consumerProperties.getString(PUBSUB_SECURITY_PROTOCOL), "SSL"); + assertEquals(consumerProperties.getString("ssl.keystore.location"), "credential-keystore"); + assertEquals(consumerProperties.getString("xc.tls.key.store.type"), "PKCS12"); + } + @Test(expectedExceptions = VeniceException.class, expectedExceptionsMessageRegExp = ".*Repush with TTL is only supported while using Kafka Input Format.*") public void testRepushTTLJobWithNonKafkaInput() { Properties repushProps = new Properties(); From 8738281041f36c33055206aa5a702516f5358e71 Mon Sep 17 00:00:00 2001 From: pthirun Date: Wed, 29 Jul 2026 11:30:41 -0700 Subject: [PATCH 2/3] [vpj] Preserve pub-sub config for dictionary training Pass the complete source consumer configuration into KafkaInputDictTrainer so Xinfra adapters, routing metadata, broker overrides, and SSL credentials are retained for both repush and empty hybrid dictionary training. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- .../linkedin/venice/hadoop/VenicePushJob.java | 11 +++- .../input/kafka/KafkaInputDictTrainer.java | 12 ++-- .../kafka/TestKafkaInputDictTrainer.java | 64 +++++++++++++++++-- 3 files changed, 73 insertions(+), 14 deletions(-) diff --git a/clients/venice-push-job/src/main/java/com/linkedin/venice/hadoop/VenicePushJob.java b/clients/venice-push-job/src/main/java/com/linkedin/venice/hadoop/VenicePushJob.java index 7e12048b853..8bb834ebb77 100755 --- a/clients/venice-push-job/src/main/java/com/linkedin/venice/hadoop/VenicePushJob.java +++ b/clients/venice-push-job/src/main/java/com/linkedin/venice/hadoop/VenicePushJob.java @@ -1663,10 +1663,14 @@ private Optional getCompressionDictionary() throws VeniceException { } private VeniceProperties getSourceDictionaryConsumerProperties() { + return getSourceDictionaryConsumerProperties(pushJobSetting.repushSourcePubsubBroker); + } + + private VeniceProperties getSourceDictionaryConsumerProperties(String sourcePubsubBroker) { return buildSourceDictionaryConsumerProperties( props, pushJobSetting.enableSSL ? sslProperties.get() : new Properties(), - pushJobSetting.repushSourcePubsubBroker); + sourcePubsubBroker); } @VisibleForTesting @@ -1686,7 +1690,6 @@ private ByteBuffer fetchOrBuildCompressionDictionary() throws VeniceException { KafkaInputDictTrainer.ParamBuilder paramBuilder = new KafkaInputDictTrainer.ParamBuilder() .setKeySchema(AvroCompatibilityHelper.toParsingForm(pushJobSetting.storeKeySchema)) .setNewKMESchemasFromController(pushJobSetting.newKmeSchemasFromController) - .setSslProperties(pushJobSetting.enableSSL ? sslProperties.get() : new Properties()) .setCompressionDictSize( props.getInt( COMPRESSION_DICTIONARY_SIZE_LIMIT, @@ -1701,6 +1704,7 @@ private ByteBuffer fetchOrBuildCompressionDictionary() throws VeniceException { LOGGER.info("Rebuild a new Zstd dictionary from the input topic: {}", pushJobSetting.kafkaInputTopic); paramBuilder.setKafkaInputBroker(pushJobSetting.repushSourcePubsubBroker) .setTopicName(pushJobSetting.kafkaInputTopic) + .setConsumerProperties(getSourceDictionaryConsumerProperties().toProperties()) .setSourceVersionCompressionStrategy(pushJobSetting.sourceKafkaInputVersionInfo.getCompressionStrategy()); KafkaInputDictTrainer dictTrainer = new KafkaInputDictTrainer(paramBuilder.build()); return ByteBuffer.wrap(dictTrainer.trainDict()); @@ -1749,8 +1753,9 @@ private ByteBuffer fetchOrBuildCompressionDictionary() throws VeniceException { "Rebuild a new Zstd dictionary from the source topic: {} in Kafka: {}", sourceTopicName, sourceKafkaUrl); - paramBuilder.setKafkaInputBroker(repushInfoResponse.getRepushInfo().getKafkaBrokerUrl()) + paramBuilder.setKafkaInputBroker(sourceKafkaUrl) .setTopicName(sourceTopicName) + .setConsumerProperties(getSourceDictionaryConsumerProperties(sourceKafkaUrl).toProperties()) .setSourceVersionCompressionStrategy( repushInfoResponse.getRepushInfo().getVersion().getCompressionStrategy()); KafkaInputDictTrainer dictTrainer = new KafkaInputDictTrainer(paramBuilder.build()); diff --git a/clients/venice-push-job/src/main/java/com/linkedin/venice/hadoop/input/kafka/KafkaInputDictTrainer.java b/clients/venice-push-job/src/main/java/com/linkedin/venice/hadoop/input/kafka/KafkaInputDictTrainer.java index b1695dcf2e2..c8ae6fb05bf 100644 --- a/clients/venice-push-job/src/main/java/com/linkedin/venice/hadoop/input/kafka/KafkaInputDictTrainer.java +++ b/clients/venice-push-job/src/main/java/com/linkedin/venice/hadoop/input/kafka/KafkaInputDictTrainer.java @@ -50,7 +50,7 @@ public static class Param { private final String kafkaInputBroker; private final String topicName; private final String keySchema; - private final Properties sslProperties; + private final Properties consumerProperties; private final int compressionDictSize; private final int dictSampleSize; private final CompressionStrategy sourceVersionCompressionStrategy; @@ -62,7 +62,7 @@ public static class Param { this.kafkaInputBroker = builder.kafkaInputBroker; this.topicName = builder.topicName; this.keySchema = builder.keySchema; - this.sslProperties = builder.sslProperties; + this.consumerProperties = builder.consumerProperties; this.compressionDictSize = builder.compressionDictSize; this.dictSampleSize = builder.dictSampleSize; this.sourceVersionCompressionStrategy = builder.sourceVersionCompressionStrategy; @@ -75,7 +75,7 @@ public static class ParamBuilder { private String kafkaInputBroker; private String topicName; private String keySchema; - private Properties sslProperties; + private Properties consumerProperties; private int compressionDictSize; private int dictSampleSize; private CompressionStrategy sourceVersionCompressionStrategy; @@ -97,8 +97,8 @@ public ParamBuilder setKeySchema(String keySchema) { return this; } - public ParamBuilder setSslProperties(Properties sslProperties) { - this.sslProperties = sslProperties; + public ParamBuilder setConsumerProperties(Properties consumerProperties) { + this.consumerProperties = consumerProperties; return this; } @@ -166,11 +166,11 @@ protected KafkaInputDictTrainer( this.trainerSupplier = trainerSupplier; this.sourceVersionCompressionStrategy = param.sourceVersionCompressionStrategy; Properties properties = new Properties(); + properties.putAll(param.consumerProperties); properties.setProperty(VENICE_REPUSH_SOURCE_PUBSUB_BROKER, param.kafkaInputBroker); properties.setProperty(KAFKA_INPUT_TOPIC, param.topicName); properties.setProperty(KAFKA_SOURCE_KEY_SCHEMA_STRING_PROP, param.keySchema); this.sourceTopicName = param.topicName; - properties.putAll(param.sslProperties); properties.setProperty(COMPRESSION_DICTIONARY_SIZE_LIMIT, Integer.toString(param.compressionDictSize)); properties.setProperty(COMPRESSION_DICTIONARY_SAMPLE_SIZE, Integer.toString(param.dictSampleSize)); properties diff --git a/clients/venice-push-job/src/test/java/com/linkedin/venice/hadoop/input/kafka/TestKafkaInputDictTrainer.java b/clients/venice-push-job/src/test/java/com/linkedin/venice/hadoop/input/kafka/TestKafkaInputDictTrainer.java index 30d8242af8d..17e10b58d11 100644 --- a/clients/venice-push-job/src/test/java/com/linkedin/venice/hadoop/input/kafka/TestKafkaInputDictTrainer.java +++ b/clients/venice-push-job/src/test/java/com/linkedin/venice/hadoop/input/kafka/TestKafkaInputDictTrainer.java @@ -1,5 +1,11 @@ package com.linkedin.venice.hadoop.input.kafka; +import static com.linkedin.venice.ConfigKeys.KAFKA_BOOTSTRAP_SERVERS; +import static com.linkedin.venice.ConfigKeys.PUBSUB_BROKER_ADDRESS; +import static com.linkedin.venice.ConfigKeys.PUBSUB_CONSUMER_ADAPTER_FACTORY_CLASS; +import static com.linkedin.venice.ConfigKeys.PUBSUB_SECURITY_PROTOCOL; +import static com.linkedin.venice.ConfigKeys.PUBSUB_TYPE_ID_TO_POSITION_CLASS_NAME_MAP; +import static com.linkedin.venice.vpj.VenicePushJobConstants.KAFKA_INPUT_TOPIC; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.doAnswer; @@ -7,6 +13,10 @@ import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; +import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertFalse; +import static org.testng.Assert.assertTrue; +import static org.testng.Assert.fail; import com.github.luben.zstd.ZstdDictTrainer; import com.linkedin.venice.compression.CompressionStrategy; @@ -37,6 +47,7 @@ import org.apache.avro.Schema; import org.apache.hadoop.mapred.InputSplit; import org.apache.hadoop.mapred.RecordReader; +import org.mockito.ArgumentCaptor; import org.testng.annotations.Test; @@ -52,6 +63,13 @@ private KafkaInputDictTrainer.Param getParam(int sampleSize) { } private KafkaInputDictTrainer.Param getParam(int sampleSize, CompressionStrategy sourceVersionCompressionStrategy) { + return getParam(sampleSize, sourceVersionCompressionStrategy, new Properties()); + } + + private KafkaInputDictTrainer.Param getParam( + int sampleSize, + CompressionStrategy sourceVersionCompressionStrategy, + Properties consumerProperties) { Map allSchemas = Utils.getAllSchemasFromResources(AvroProtocolDefinition.KAFKA_MESSAGE_ENVELOPE); Map allSchemaStr = allSchemas.entrySet().stream().collect(Collectors.toMap(Map.Entry::getKey, e -> e.getValue().toString())); @@ -60,14 +78,14 @@ private KafkaInputDictTrainer.Param getParam(int sampleSize, CompressionStrategy .setKeySchema("\"string\"") .setCompressionDictSize(900 * 1024) .setDictSampleSize(sampleSize) - .setSslProperties(new Properties()) + .setConsumerProperties(consumerProperties) .setSourceVersionCompressionStrategy(sourceVersionCompressionStrategy) .setNewKMESchemasFromController(allSchemaStr) .build(); } - @Test(expectedExceptions = VeniceException.class, expectedExceptionsMessageRegExp = "No record.*") - public void testEmptyTopic() throws IOException { + @Test + public void testConsumerPropertiesPropagateToEmptyTopicTraining() throws IOException { KafkaInputFormat mockFormat = mock(KafkaInputFormat.class); PubSubTopicPartition topicPartition = new PubSubTopicPartitionImpl(PUB_SUB_TOPIC_REPOSITORY.getTopic("test_topic"), 0); @@ -79,12 +97,48 @@ public void testEmptyTopic() throws IOException { doReturn(false).when(mockRecordReader).next(any(), any()); doReturn(mockRecordReader).when(mockFormat).getRecordReader(any(), any(), any(), any()); + Properties consumerProperties = new Properties(); + consumerProperties.setProperty( + PUBSUB_CONSUMER_ADAPTER_FACTORY_CLASS, + "com.linkedin.venice.pubsub.adapter.xinfra.consumer.XcConsumerAdapterFactory"); + consumerProperties.setProperty( + PUBSUB_TYPE_ID_TO_POSITION_CLASS_NAME_MAP, + "1:com.linkedin.venice.pubsub.adapter.xinfra.XinfraPosition"); + consumerProperties.setProperty("xc.pubsub.broker.url.to.region.name.map", "northguard:ei4"); + consumerProperties.setProperty(PUBSUB_BROKER_ADDRESS, "test_url"); + consumerProperties.setProperty(KAFKA_BOOTSTRAP_SERVERS, "test_url"); + consumerProperties.setProperty(PUBSUB_SECURITY_PROTOCOL, "SSL"); + consumerProperties.setProperty("ssl.keystore.location", "credential-keystore"); + KafkaInputDictTrainer trainer = new KafkaInputDictTrainer( mockFormat, Optional.empty(), - getParam(100), + getParam(100, CompressionStrategy.NO_OP, consumerProperties), getCompressorBuilder(new NoopCompressor())); - trainer.trainDict(Optional.of(mock(PubSubConsumerAdapter.class))); + try { + trainer.trainDict(Optional.of(mock(PubSubConsumerAdapter.class))); + fail("Expected training on an empty topic to fail"); + } catch (VeniceException e) { + assertTrue(e.getMessage().startsWith("No record")); + } + + ArgumentCaptor consumerPropertiesCaptor = ArgumentCaptor.forClass(VeniceProperties.class); + verify(mockFormat).getSplits(consumerPropertiesCaptor.capture()); + VeniceProperties actualConsumerProperties = consumerPropertiesCaptor.getValue(); + assertEquals( + actualConsumerProperties.getString(PUBSUB_CONSUMER_ADAPTER_FACTORY_CLASS), + "com.linkedin.venice.pubsub.adapter.xinfra.consumer.XcConsumerAdapterFactory"); + assertEquals( + actualConsumerProperties.getString(PUBSUB_TYPE_ID_TO_POSITION_CLASS_NAME_MAP), + "1:com.linkedin.venice.pubsub.adapter.xinfra.XinfraPosition"); + assertEquals(actualConsumerProperties.getString("xc.pubsub.broker.url.to.region.name.map"), "northguard:ei4"); + assertEquals(actualConsumerProperties.getString(PUBSUB_BROKER_ADDRESS), "test_url"); + assertEquals(actualConsumerProperties.getString(KAFKA_BOOTSTRAP_SERVERS), "test_url"); + assertEquals(actualConsumerProperties.getString(PUBSUB_SECURITY_PROTOCOL), "SSL"); + assertEquals(actualConsumerProperties.getString("ssl.keystore.location"), "credential-keystore"); + assertFalse( + consumerProperties.containsKey(KAFKA_INPUT_TOPIC), + "Building trainer properties must not mutate the supplied consumer properties"); } interface ResettableRecordReader extends RecordReader { From ca71159abadb8941e0f0ba6019bb9858249ce0b2 Mon Sep 17 00:00:00 2001 From: pthirun Date: Wed, 29 Jul 2026 11:44:13 -0700 Subject: [PATCH 3/3] [vpj] Cover dictionary consumer SSL branch Exercise the production dictionary consumer property helper without SSL so the new conditional path is represented in VPJ diff coverage. This keeps the propagation assertions on the same helper used by dictionary reads and training. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- .../linkedin/venice/hadoop/VenicePushJob.java | 3 +- .../hadoop/VenicePushJobRepushTest.java | 33 +++++++++---------- 2 files changed, 18 insertions(+), 18 deletions(-) diff --git a/clients/venice-push-job/src/main/java/com/linkedin/venice/hadoop/VenicePushJob.java b/clients/venice-push-job/src/main/java/com/linkedin/venice/hadoop/VenicePushJob.java index 8bb834ebb77..702b75c976d 100755 --- a/clients/venice-push-job/src/main/java/com/linkedin/venice/hadoop/VenicePushJob.java +++ b/clients/venice-push-job/src/main/java/com/linkedin/venice/hadoop/VenicePushJob.java @@ -1666,7 +1666,8 @@ private VeniceProperties getSourceDictionaryConsumerProperties() { return getSourceDictionaryConsumerProperties(pushJobSetting.repushSourcePubsubBroker); } - private VeniceProperties getSourceDictionaryConsumerProperties(String sourcePubsubBroker) { + @VisibleForTesting + VeniceProperties getSourceDictionaryConsumerProperties(String sourcePubsubBroker) { return buildSourceDictionaryConsumerProperties( props, pushJobSetting.enableSSL ? sslProperties.get() : new Properties(), diff --git a/clients/venice-push-job/src/test/java/com/linkedin/venice/hadoop/VenicePushJobRepushTest.java b/clients/venice-push-job/src/test/java/com/linkedin/venice/hadoop/VenicePushJobRepushTest.java index 725d308c7cc..d0194d750b6 100644 --- a/clients/venice-push-job/src/test/java/com/linkedin/venice/hadoop/VenicePushJobRepushTest.java +++ b/clients/venice-push-job/src/test/java/com/linkedin/venice/hadoop/VenicePushJobRepushTest.java @@ -61,24 +61,23 @@ public void testSourceDictionaryConsumerPropertiesRetainPubSubConfigAndOverrideB jobProperties.setProperty(PUBSUB_BROKER_ADDRESS, "destination-broker"); jobProperties.setProperty(KAFKA_BOOTSTRAP_SERVERS, "legacy-broker"); - VeniceProperties consumerProperties = VenicePushJob.buildSourceDictionaryConsumerProperties( - new VeniceProperties(jobProperties), - new Properties(), - "source-broker"); + try (VenicePushJob pushJob = getSpyVenicePushJob(jobProperties, null)) { + VeniceProperties consumerProperties = pushJob.getSourceDictionaryConsumerProperties("source-broker"); - assertEquals( - consumerProperties.getString(PUBSUB_CONSUMER_ADAPTER_FACTORY_CLASS), - "com.linkedin.venice.pubsub.adapter.xinfra.consumer.XcConsumerAdapterFactory"); - assertEquals( - consumerProperties.getString(PUBSUB_TYPE_ID_TO_POSITION_CLASS_NAME_MAP), - "1:com.linkedin.venice.pubsub.adapter.xinfra.XinfraPosition"); - assertEquals(consumerProperties.getString("xc.pubsub.broker.url.to.region.name.map"), "northguard:ei4"); - assertEquals(consumerProperties.getString(PUBSUB_BROKER_ADDRESS), "source-broker"); - assertEquals(consumerProperties.getString(KAFKA_BOOTSTRAP_SERVERS), "source-broker"); - assertEquals( - jobProperties.getProperty(PUBSUB_BROKER_ADDRESS), - "destination-broker", - "Building consumer properties must not mutate the job properties"); + assertEquals( + consumerProperties.getString(PUBSUB_CONSUMER_ADAPTER_FACTORY_CLASS), + "com.linkedin.venice.pubsub.adapter.xinfra.consumer.XcConsumerAdapterFactory"); + assertEquals( + consumerProperties.getString(PUBSUB_TYPE_ID_TO_POSITION_CLASS_NAME_MAP), + "1:com.linkedin.venice.pubsub.adapter.xinfra.XinfraPosition"); + assertEquals(consumerProperties.getString("xc.pubsub.broker.url.to.region.name.map"), "northguard:ei4"); + assertEquals(consumerProperties.getString(PUBSUB_BROKER_ADDRESS), "source-broker"); + assertEquals(consumerProperties.getString(KAFKA_BOOTSTRAP_SERVERS), "source-broker"); + assertEquals( + jobProperties.getProperty(PUBSUB_BROKER_ADDRESS), + "destination-broker", + "Building consumer properties must not mutate the job properties"); + } } @Test