Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -2818,6 +2818,26 @@ private ConfigKeys() {
*/
public static final String WRITER_BATCHING_MAX_BUFFER_SIZE_IN_BYTES = "writer.batching.max.buffer.size.in.bytes";

/**
* Controls when {@link com.linkedin.venice.writer.VeniceWriter} attaches the
* {@link com.linkedin.venice.pubsub.api.PubSubMessageHeaders#VENICE_TRANSPORT_PROTOCOL_HEADER vtp}
* protocol-schema header to outbound messages. The pre-existing emission gate is
* {@code segmentNumber == 0 && messageSequenceNumber == 0} on the outgoing producer metadata.
* On the data path that gate matches only the first segment-start record produced on a
* partition (segment 0, sequence 0); subsequent data SOS records use non-zero
* {@code segmentNumber} and are unaffected. Heartbeats pin both coordinates to {@code 0}, so
* every heartbeat matches the gate. Accepts one of the
* {@link com.linkedin.venice.writer.VtpHeaderEmissionMode} names: {@code SOS_AND_HB} (default,
* preserves the pre-existing behavior — emit on both first data SOS and every heartbeat),
* {@code SOS_ONLY} (apply the same 0/0 gate but skip heartbeats — i.e., emit on the first
* data SOS per partition and on DoL stamps, which also carry 0/0 coordinates but are not
* heartbeats), or {@code NONE} (never emit). Use {@code SOS_ONLY} when heartbeat
* fan-out dominates the consumer-side per-record memory footprint and consumers can bootstrap
* the {@code KafkaMessageEnvelope} schema from the first data SOS or an out-of-band schema
* cache.
*/
public static final String VENICE_WRITER_VTP_HEADER_EMISSION_MODE = "venice.writer.vtp.header.emission.mode";

/**
* The maximum age (in milliseconds) of producer state retained by Data Ingestion Validation. Tuning this
* can prevent OOMing in cases where there is a lot of historical churn in RT producers. The age of a given
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

import static com.linkedin.venice.ConfigKeys.INSTANCE_ID;
import static com.linkedin.venice.ConfigKeys.LISTENER_PORT;
import static com.linkedin.venice.ConfigKeys.VENICE_WRITER_VTP_HEADER_EMISSION_MODE;
import static com.linkedin.venice.message.KafkaKey.CONTROL_MESSAGE_KAFKA_KEY_LENGTH;
import static com.linkedin.venice.pubsub.api.PubSubMessageHeaders.VENICE_TRANSPORT_PROTOCOL_HEADER;
import static com.linkedin.venice.pubsub.api.PubSubMessageHeaders.VENICE_VIEW_PARTITIONS_MAP_HEADER;
Expand Down Expand Up @@ -252,6 +253,7 @@ public class VeniceWriter<K, V, U> extends AbstractVeniceWriter<K, V, U> {

// Immutable state
private final PubSubMessageHeader protocolSchemaHeader;
private final VtpHeaderEmissionMode vtpHeaderEmissionMode;

protected final VeniceKafkaSerializer<K> keySerializer;
protected final VeniceKafkaSerializer<V> valueSerializer;
Expand Down Expand Up @@ -439,6 +441,29 @@ public VeniceWriter(
: new PubSubMessageHeader(
VENICE_TRANSPORT_PROTOCOL_HEADER,
overrideProtocolSchema.toString().getBytes(StandardCharsets.UTF_8));
/*
* Parse VENICE_WRITER_VTP_HEADER_EMISSION_MODE. Default is SOS_AND_HB to preserve the
* pre-existing emission rule (attach vtp when segmentNumber == 0 && messageSequenceNumber == 0):
* on the data path that gate matches only the first segment-start record per partition (segment
* 0, sequence 0), and every heartbeat (heartbeats pin both coordinates to 0 via
* getHeartbeatKME(...)). Unknown values fall back to the default with a warning.
*/
String vtpHeaderEmissionModeProp =
props.getString(VENICE_WRITER_VTP_HEADER_EMISSION_MODE, VtpHeaderEmissionMode.SOS_AND_HB.name());
VtpHeaderEmissionMode parsedMode;
try {
parsedMode = vtpHeaderEmissionModeProp == null
? VtpHeaderEmissionMode.SOS_AND_HB
: VtpHeaderEmissionMode.valueOf(vtpHeaderEmissionModeProp.trim());
} catch (IllegalArgumentException e) {
logger.warn(
"Unrecognized {} value '{}'; falling back to {}",
VENICE_WRITER_VTP_HEADER_EMISSION_MODE,
vtpHeaderEmissionModeProp,
VtpHeaderEmissionMode.SOS_AND_HB);
parsedMode = VtpHeaderEmissionMode.SOS_AND_HB;
}
this.vtpHeaderEmissionMode = parsedMode;

try {
this.producerAdapter = producerAdapter;
Expand Down Expand Up @@ -1882,6 +1907,7 @@ private CompletableFuture<PubSubProduceResult> sendMessage(
PubSubProducerCallback outputCallback = setInternalCallback(callback, internalCallback);
PubSubMessageHeaders finalPubSubMessageHeaders = getHeaders(
kafkaValue.getProducerMetadata(),
false /* isHeartbeat */,
false,
LeaderCompleteState.LEADER_NOT_COMPLETED,
pubSubMessageHeaders);
Expand Down Expand Up @@ -1913,8 +1939,11 @@ PubSubProducerCallback setInternalCallback(
}

/**
* {@link PubSubMessageHeaders#VENICE_TRANSPORT_PROTOCOL_HEADER} or {@link EmptyPubSubMessageHeaders} is used for
* all messages to a partition based on {@link VeniceWriter} param overrideProtocolSchema and whether it's a first message.
* Builds the {@link PubSubMessageHeaders} for an outbound message. The vtp protocol-schema header
* ({@link PubSubMessageHeaders#VENICE_TRANSPORT_PROTOCOL_HEADER}) is attached when all of: (1) overrideProtocolSchema
* is non-null, (2) the message is segment 0 / sequence 0, and (3) {@link VtpHeaderEmissionMode} permits emission for
* this message type (heartbeat vs non-heartbeat). Under {@code SOS_ONLY} heartbeat SOS records are skipped; under
* {@code NONE} no message gets the header regardless.
* {@link PubSubMessageHeaders#VENICE_LEADER_COMPLETION_STATE_HEADER} is added to the above headers for HB SOS message.
* {@link PubSubMessageHeaders#VENICE_VIEW_PARTITIONS_MAP_HEADER} is added to the headers for chunked messages
* of materialized
Expand All @@ -1931,13 +1960,36 @@ PubSubProducerCallback setInternalCallback(

private PubSubMessageHeaders getHeaders(
ProducerMetadata producerMetadata,
boolean isHeartbeat,
boolean addLeaderCompleteState,
LeaderCompleteState leaderCompleteState,
PubSubMessageHeaders headers) {
Comment thread
sushantmane marked this conversation as resolved.
PubSubMessageHeader viewPartitionHeader = headers.get(VENICE_VIEW_PARTITIONS_MAP_HEADER);
// If the message is the first message in a segment, we need to add the protocol schema headers.
boolean needVtpHeader =
/*
* Decide whether to attach the vtp protocol-schema header on this outbound message.
*
* Pre-existing rule: attach on the first message of the first segment, i.e. SOS records
* (segmentNumber == 0 && messageSequenceNumber == 0). Heartbeats are encoded as
* START_OF_SEGMENT with both numbers zero, so under SOS_AND_HB every heartbeat picks up
* the ~16 KB vtp blob — that dominates the per-record memory footprint on busy ingestion
* paths and is the lever VtpHeaderEmissionMode lets writers tune. See
* VtpHeaderEmissionMode for the per-mode semantics.
*/
boolean isFirstMessageOfFirstSegment =
producerMetadata.getSegmentNumber() == 0 && producerMetadata.getMessageSequenceNumber() == 0;
boolean needVtpHeader;
switch (vtpHeaderEmissionMode) {
case NONE:
needVtpHeader = false;
break;
case SOS_ONLY:
needVtpHeader = isFirstMessageOfFirstSegment && !isHeartbeat;
break;
case SOS_AND_HB:
default:
needVtpHeader = isFirstMessageOfFirstSegment;
break;
}

// construct PubSubMessageHeaders only if it is needed
PubSubMessageHeaders returnPubSubMessageHeaders = (headers instanceof EmptyPubSubMessageHeaders)
Expand Down Expand Up @@ -2547,7 +2599,12 @@ public CompletableFuture<PubSubProduceResult> sendDoLStamp(
* lands on the wire with empty headers, and a forward-compat consumer that hits it as
* the first record on a fresh VT has no way to bootstrap an unknown KME schema.
*/
getHeaders(kafkaMessageEnvelope.getProducerMetadata(), false, null, EmptyPubSubMessageHeaders.SINGLETON),
getHeaders(
kafkaMessageEnvelope.getProducerMetadata(),
false /* isHeartbeat */,
false,
null,
EmptyPubSubMessageHeaders.SINGLETON),
callback);
}
}
Expand Down Expand Up @@ -2601,6 +2658,7 @@ public CompletableFuture<PubSubProduceResult> sendHeartbeat(
kafkaMessageEnvelope,
getHeaders(
kafkaMessageEnvelope.getProducerMetadata(),
true /* isHeartbeat */,
addLeaderCompleteState,
leaderCompleteState,
EmptyPubSubMessageHeaders.SINGLETON),
Expand All @@ -2624,6 +2682,7 @@ public CompletableFuture<PubSubProduceResult> sendHeartbeat(
kafkaMessageEnvelope,
getHeaders(
kafkaMessageEnvelope.getProducerMetadata(),
true /* isHeartbeat */,
addLeaderCompleteState,
leaderCompleteState,
EmptyPubSubMessageHeaders.SINGLETON),
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
package com.linkedin.venice.writer;

/**
* Controls whether {@link VeniceWriter} attaches the
* {@link com.linkedin.venice.pubsub.api.PubSubMessageHeaders#VENICE_TRANSPORT_PROTOCOL_HEADER vtp}
* protocol-schema header to outbound messages.
*
* <p>The vtp header carries the Avro JSON for {@code KafkaMessageEnvelope} (~16 KB) and is only
* useful to consumers that need to bootstrap the schema for forward compatibility. Pre-existing
* behavior emits it on outbound messages whose producer metadata satisfies
* {@code segmentNumber == 0 && messageSequenceNumber == 0}. On the data path that gate only
* matches the very first segment-start record produced on a partition (segment 0, sequence 0);
* subsequent data SOS records use non-zero {@code segmentNumber} and are not affected. Heartbeat
* control messages, however, are synthesized with both coordinates pinned to {@code 0} (see
* {@link VeniceWriter#getHeartbeatKME}), so every heartbeat matches the gate and picks up the
* ~16 KB header even though heartbeat consumers never use it for schema bootstrap. On busy
* ingestion paths with many partitions and frequent heartbeats this dominates the per-record
* memory footprint.
*
* <p>This mode lets writers opt out of the heartbeat case (or out entirely) without changing the
* semantics of regular data segment-start records.
*/
public enum VtpHeaderEmissionMode {
/**
* Emit the vtp header on every segment-start message, including heartbeat SOS records. Default,
* preserves pre-existing behavior.
*/
SOS_AND_HB,

/**
* Emit the vtp header on non-heartbeat segment-start messages (data SOS and DoL stamps); skip
* heartbeat SOS records. DoL stamps are START_OF_SEGMENT control messages that carry
* segmentNumber=0 / messageSequenceNumber=0 (like heartbeats) but are not heartbeats — they
* need the vtp header so a forward-compat consumer hitting a DoL as the first record on a
* fresh VT can bootstrap the KME schema. Consumers are expected to obtain the
* {@code KafkaMessageEnvelope} schema for heartbeat records by some other means (controller-side
* schema cache, classpath fallback, or an earlier data SOS on the same segment).
*/
SOS_ONLY,

/**
* Never emit the vtp header. Use only when all consumers can resolve the
* {@code KafkaMessageEnvelope} schema without the per-segment hint.
*/
NONE
}
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
package com.linkedin.venice.writer;

import static com.linkedin.venice.ConfigKeys.VENICE_WRITER_VTP_HEADER_EMISSION_MODE;
import static com.linkedin.venice.message.KafkaKey.DOL_STAMP;
import static com.linkedin.venice.message.KafkaKey.HEART_BEAT;
import static com.linkedin.venice.pubsub.api.PubSubMessageHeaders.EXECUTION_ID_KEY;
Expand Down Expand Up @@ -57,6 +58,7 @@
import com.linkedin.venice.pubsub.api.PubSubMessageHeader;
import com.linkedin.venice.pubsub.api.PubSubMessageHeaders;
import com.linkedin.venice.pubsub.api.PubSubPosition;
import com.linkedin.venice.pubsub.api.PubSubProduceResult;
import com.linkedin.venice.pubsub.api.PubSubProducerAdapter;
import com.linkedin.venice.pubsub.api.PubSubProducerCallback;
import com.linkedin.venice.pubsub.api.PubSubTopic;
Expand Down Expand Up @@ -1845,6 +1847,116 @@ public void testPassThroughPutHeaderHandling(
}
}

@DataProvider(name = "VtpHeaderEmissionMode-HeartbeatExpectsVtp")
public static Object[][] vtpHeaderEmissionModeHeartbeatExpectsVtp() {
return new Object[][] { { VtpHeaderEmissionMode.SOS_AND_HB, true }, { VtpHeaderEmissionMode.SOS_ONLY, false },
{ VtpHeaderEmissionMode.NONE, false } };
}

/**
* Verifies that {@link VeniceWriter#sendHeartbeat} respects {@link VtpHeaderEmissionMode}: the
* {@code vtp} protocol-schema header is attached only when the configured mode is
* {@link VtpHeaderEmissionMode#SOS_AND_HB}. Default is {@code SOS_AND_HB} so existing callers
* are unaffected; {@code SOS_ONLY} and {@code NONE} both drop the header on heartbeats.
*/
@Test(dataProvider = "VtpHeaderEmissionMode-HeartbeatExpectsVtp")
public void testHeartbeatVtpEmissionMode(VtpHeaderEmissionMode mode, boolean expectVtpHeader) {
PubSubProducerAdapter mockedProducer = mock(PubSubProducerAdapter.class);
CompletableFuture mockedFuture = mock(CompletableFuture.class);
when(mockedProducer.sendMessage(any(), any(), any(), any(), any(), any())).thenReturn(mockedFuture);

Properties writerProperties = new Properties();
writerProperties.setProperty(VENICE_WRITER_VTP_HEADER_EMISSION_MODE, mode.name());

String stringSchema = "\"string\"";
VeniceKafkaSerializer serializer = new VeniceAvroKafkaSerializer(stringSchema);
String testTopic = "test_rt_vtp_mode";
VeniceWriterOptions opts = new VeniceWriterOptions.Builder(testTopic).setKeyPayloadSerializer(serializer)
.setValuePayloadSerializer(serializer)
.setWriteComputePayloadSerializer(serializer)
.setPartitioner(new DefaultVenicePartitioner())
.setTime(SystemTime.INSTANCE)
.setPartitionCount(1)
.build();
VeniceWriter<Object, Object, Object> writer =
new VeniceWriter(opts, new VeniceProperties(writerProperties), mockedProducer);

PubSubTopicPartition topicPartition = mock(PubSubTopicPartition.class);
PubSubTopic topic = mock(PubSubTopic.class);
when(topic.getName()).thenReturn(testTopic);
when(topicPartition.getPubSubTopic()).thenReturn(topic);
when(topicPartition.getPartitionNumber()).thenReturn(0);

writer.sendHeartbeat(
topicPartition,
null,
DEFAULT_LEADER_METADATA_WRAPPER,
false /* addLeaderCompleteState */,
LEADER_NOT_COMPLETED,
System.currentTimeMillis());

ArgumentCaptor<PubSubMessageHeaders> headersCaptor = ArgumentCaptor.forClass(PubSubMessageHeaders.class);
ArgumentCaptor<KafkaKey> keyCaptor = ArgumentCaptor.forClass(KafkaKey.class);
verify(mockedProducer, times(1))
.sendMessage(eq(testTopic), eq(0), keyCaptor.capture(), any(), headersCaptor.capture(), any());

/* Sanity-check that the message we captured really is a heartbeat — uses the HEART_BEAT key. */
assertTrue(Arrays.equals(HEART_BEAT.getKey(), keyCaptor.getValue().getKey()));

PubSubMessageHeaders capturedHeaders = headersCaptor.getValue();
boolean vtpAttached = capturedHeaders.get(VENICE_TRANSPORT_PROTOCOL_HEADER) != null;
assertEquals(
vtpAttached,
expectVtpHeader,
"VtpHeaderEmissionMode=" + mode + ": vtp header on heartbeat should be "
+ (expectVtpHeader ? "present" : "absent") + " but was " + (vtpAttached ? "present" : "absent"));
}

/**
* Sanity-checks that {@link VtpHeaderEmissionMode#NONE} also drops the {@code vtp} header on
* regular data segment-start records (not just heartbeats). The first {@link VeniceWriter#put}
* call on a fresh writer produces the segment-start message, which is where {@code vtp} would
* be attached under {@code SOS_AND_HB} or {@code SOS_ONLY}.
*/
@Test
public void testDataSosWithVtpEmissionModeNone() throws ExecutionException, InterruptedException {
PubSubProducerAdapter mockedProducer = mock(PubSubProducerAdapter.class);
CompletableFuture<PubSubProduceResult> done = CompletableFuture.completedFuture(null);
when(mockedProducer.sendMessage(any(), any(), any(), any(), any(), any())).thenReturn(done);

Properties writerProperties = new Properties();
writerProperties.setProperty(VENICE_WRITER_VTP_HEADER_EMISSION_MODE, VtpHeaderEmissionMode.NONE.name());

String stringSchema = "\"string\"";
VeniceKafkaSerializer serializer = new VeniceAvroKafkaSerializer(stringSchema);
String testTopic = "test_data_vtp_none";
VeniceWriterOptions opts = new VeniceWriterOptions.Builder(testTopic).setKeyPayloadSerializer(serializer)
.setValuePayloadSerializer(serializer)
.setWriteComputePayloadSerializer(serializer)
.setPartitioner(new DefaultVenicePartitioner())
.setTime(SystemTime.INSTANCE)
.setPartitionCount(1)
.build();
VeniceWriter<Object, Object, Object> writer =
new VeniceWriter(opts, new VeniceProperties(writerProperties), mockedProducer);
writer.put("k0", "v0", 1, null);

ArgumentCaptor<PubSubMessageHeaders> headersCaptor = ArgumentCaptor.forClass(PubSubMessageHeaders.class);
verify(mockedProducer, atLeast(1))
.sendMessage(eq(testTopic), anyInt(), any(), any(), headersCaptor.capture(), any());

/*
* Under NONE the vtp header must be absent on every emitted message, including the first
* data record (which under SOS_AND_HB would carry it because it is the first message of the
* first segment).
*/
for (PubSubMessageHeaders headers: headersCaptor.getAllValues()) {
assertNull(
headers.get(VENICE_TRANSPORT_PROTOCOL_HEADER),
"VtpHeaderEmissionMode=NONE must not attach a vtp header to any emitted message");
}
}

private static KafkaMessageEnvelope buildEopEnvelope() {
KafkaMessageEnvelope kme = new KafkaMessageEnvelope();
kme.messageType = MessageType.CONTROL_MESSAGE.getValue();
Expand Down