Skip to content
Open
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 @@ -138,6 +138,7 @@ public boolean isReplicaLaggedAndNeedBlobTransfer(
LOGGER.warn("Offset record not found for: {}", replicaId);
return true;
}
offsetRecord.setReplicaId(replicaId);
// TODO: Remove offset lag threshold entirely in the future — no offset lag should be allowed.
long blobTransferDisabledOffsetLagThreshold = serverConfig.getBlobTransferDisabledOffsetLagThreshold();
if (blobTransferDisabledOffsetLagThreshold < 0) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -426,6 +426,7 @@ public PartitionConsumptionState(
}
this.hybrid = hybrid;
this.offsetRecord = offsetRecord;
this.offsetRecord.setReplicaId(getReplicaId());
this.pubSubContext = pubSubContext;
this.errorReported = false;
this.lagCaughtUp = false;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6103,7 +6103,8 @@ protected PubSubPosition extractUpstreamPosition(DefaultPubSubMessage consumerRe
return PubSubUtil.deserializePositionWithOffsetFallback(
leaderMetadataFooter.upstreamPubSubPosition,
leaderMetadataFooter.upstreamOffset,
pubSubContext.getPubSubPositionDeserializer());
pubSubContext.getPubSubPositionDeserializer(),
getReplicaId(kafkaVersionTopic, consumerRecord.getPartition()));
} else {
// Directly use upstreamOffset without attempting position deserialization
return PubSubUtil.fromKafkaOffset(leaderMetadataFooter.upstreamOffset);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
import com.linkedin.venice.utils.CollectionUtils;
import com.linkedin.venice.utils.LatencyUtils;
import com.linkedin.venice.utils.RedundantExceptionFilter;
import com.linkedin.venice.utils.Utils;
import com.linkedin.venice.utils.concurrent.VeniceConcurrentHashMap;
import com.linkedin.venice.utils.lazy.Lazy;
import java.nio.ByteBuffer;
Expand Down Expand Up @@ -1003,7 +1004,8 @@ private String generateMessage(
deserializePositionWithOffsetFallback(
consumerRecord.getValue().leaderMetadataFooter.upstreamPubSubPosition,
consumerRecord.getValue().leaderMetadataFooter.upstreamOffset,
pubSubPositionDeserializer))
pubSubPositionDeserializer,
Utils.getReplicaId(consumerRecord.getTopicPartition())))
.append("; upstream pub sub cluster ID: ")
.append(consumerRecord.getValue().leaderMetadataFooter.upstreamKafkaClusterId)
.append("; producer host name: ")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -503,7 +503,8 @@ private void logRecordMetadata(DefaultPubSubMessage record) {
: PubSubUtil.deserializePositionWithOffsetFallback(
leaderMetadata.upstreamPubSubPosition,
leaderMetadata.upstreamOffset,
pubSubPositionDeserializer),
pubSubPositionDeserializer,
Utils.getReplicaId(record.getTopicPartition())),
leaderMetadata == null ? "-" : leaderMetadata.upstreamKafkaClusterId,
chunkMetadata);
} catch (Exception e) {
Expand Down Expand Up @@ -669,7 +670,8 @@ static String constructTopicSwitchLog(
: PubSubUtil.deserializePositionWithOffsetFallback(
leaderMetadata.upstreamPubSubPosition,
leaderMetadata.upstreamOffset,
pubSubPositionDeserializer),
pubSubPositionDeserializer,
Utils.getReplicaId(record.getTopicPartition())),
leaderMetadata == null ? "-" : leaderMetadata.upstreamKafkaClusterId);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
import com.linkedin.venice.pubsub.api.PubSubPosition;
import com.linkedin.venice.pubsub.api.PubSubTopicPartition;
import com.linkedin.venice.utils.ByteUtils;
import com.linkedin.venice.utils.Utils;
import com.linkedin.venice.vpj.pubsub.input.PubSubSplitIterator;
import java.io.IOException;
import java.util.ArrayList;
Expand Down Expand Up @@ -80,7 +81,8 @@ public static LilyPadUtils.Snapshot<ComparablePubSubPosition> buildSnapshot(
PubSubPosition rawPosition = PubSubUtil.deserializePositionWithOffsetFallback(
leaderMetadata.upstreamPubSubPosition,
leaderMetadata.upstreamOffset,
pubSubPositionDeserializer);
pubSubPositionDeserializer,
Utils.getReplicaId(topicPartition));
ComparablePubSubPosition upstreamPosition = new ComparablePubSubPosition(rawPosition, consumer, topicPartition);
ComparablePubSubPosition currentHw = runningHighWatermark.get(regionId);
if (currentHw == null || upstreamPosition.compareTo(currentHw) > 0) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,11 @@ public class OffsetRecord {
private final PubSubContext pubSubContext;
private final PubSubPositionDeserializer pubSubPositionDeserializer;

/**
* Runtime-only diagnostic context. This is intentionally kept outside {@link PartitionState}, so binding an
* OffsetRecord to a replica does not change its serialized form.
*/
private volatile String replicaId;
private PubSubTopic leaderPubSubTopic;

public OffsetRecord(
Expand Down Expand Up @@ -136,7 +141,8 @@ public void clearPreviousStatusesEntry(CharSequence key) {
public PubSubPosition getCheckpointedLocalVtPosition() {
return deserializePositionWithOffsetFallback(
this.partitionState.lastProcessedVersionTopicPubSubPosition,
this.partitionState.offset);
this.partitionState.offset,
positionSourceOrDefault("OffsetRecord.lastProcessedVersionTopicPubSubPosition"));
}

public void checkpointLocalVtPosition(PubSubPosition vtPosition) {
Expand All @@ -151,7 +157,8 @@ public void checkpointLocalVtPosition(PubSubPosition vtPosition) {
public PubSubPosition getCheckpointedRemoteVtPosition() {
return deserializePositionWithOffsetFallback(
this.partitionState.upstreamVersionTopicPubSubPosition,
this.partitionState.upstreamVersionTopicOffset);
this.partitionState.upstreamVersionTopicOffset,
positionSourceOrDefault("OffsetRecord.upstreamVersionTopicPubSubPosition"));
}

public void checkpointRemoteVtPosition(PubSubPosition remoteVtPosition) {
Expand Down Expand Up @@ -311,7 +318,10 @@ public PubSubPosition getCheckpointedRtPosition(String pubSubBrokerAddress) {
// If the offset is not set, return EARLIEST symbolic position.
return PubSubSymbolicPosition.EARLIEST;
}
return deserializePositionWithOffsetFallback(wfBuffer, offset);
return deserializePositionWithOffsetFallback(
wfBuffer,
offset,
positionSourceOrDefault("OffsetRecord.upstreamRealTimeTopicPubSubPosition[" + pubSubBrokerAddress + "]"));
}

public void checkpointRtPosition(String pubSubBrokerAddress, PubSubPosition leaderPosition) {
Expand Down Expand Up @@ -341,12 +351,36 @@ public void cloneRtPositionCheckpoints(@Nonnull Map<String, PubSubPosition> chec
for (Map.Entry<String, Long> offsetEntry: partitionState.upstreamOffsetMap.entrySet()) {
String pubSubBrokerAddress = offsetEntry.getKey();
ByteBuffer wfBuffer = partitionState.upstreamRealTimeTopicPubSubPositionMap.get(pubSubBrokerAddress);
checkpointUpstreamPositionsReceiver
.put(pubSubBrokerAddress, deserializePositionWithOffsetFallback(wfBuffer, offsetEntry.getValue()));
checkpointUpstreamPositionsReceiver.put(
pubSubBrokerAddress,
deserializePositionWithOffsetFallback(
wfBuffer,
offsetEntry.getValue(),
positionSourceOrDefault(
"OffsetRecord.upstreamRealTimeTopicPubSubPosition[" + pubSubBrokerAddress + "]")));
}
}
}

/**
* Associates runtime diagnostic context with this record. Rebinding to the same replica is harmless, while
* rebinding to a different replica is rejected because an OffsetRecord must not be shared across partitions.
*
* <p>The synchronized write and volatile read make the binding safely visible to threads that subsequently
* deserialize checkpointed positions.</p>
*/
public synchronized void setReplicaId(String replicaId) {
if (replicaId == null || replicaId.isEmpty()) {
throw new IllegalArgumentException("Replica ID cannot be null or empty");
}
if (this.replicaId != null && !this.replicaId.equals(replicaId)) {
throw new IllegalStateException(
"OffsetRecord is already associated with replica " + this.replicaId + " and cannot be rebound to "
+ replicaId);
}
this.replicaId = replicaId;
}

public GUID getLeaderGUID() {
return this.partitionState.leaderGUID;
}
Expand Down Expand Up @@ -551,10 +585,29 @@ public byte[] toBytes() {
*/
@VisibleForTesting
PubSubPosition deserializePositionWithOffsetFallback(ByteBuffer wireFormatBytes, long offset) {
return deserializePositionWithOffsetFallback(wireFormatBytes, offset, "OffsetRecord");
}

/**
* Same as {@link #deserializePositionWithOffsetFallback(ByteBuffer, long)}, but forwards an explicit
* {@code positionSource} label down to {@link PubSubUtil#deserializePositionWithOffsetFallback} so that a
* deserialization-failure warning can be attributed to the checkpoint field (and replica, when known) it came
* from, instead of logging a bare "N/A".
*/
private PubSubPosition deserializePositionWithOffsetFallback(
ByteBuffer wireFormatBytes,
long offset,
String positionSource) {
if (pubSubContext != null && !pubSubContext.isUseCheckpointedPubSubPositionWithFallbackEnabled()) {
// When feature flag is disabled, use offset-only approach without attempting wire format deserialization
return PubSubUtil.fromKafkaOffset(offset);
}
return PubSubUtil.deserializePositionWithOffsetFallback(wireFormatBytes, offset, pubSubPositionDeserializer);
return PubSubUtil
.deserializePositionWithOffsetFallback(wireFormatBytes, offset, pubSubPositionDeserializer, positionSource);
}

private String positionSourceOrDefault(String fallbackLabel) {
String currentReplicaId = replicaId;
return currentReplicaId != null ? currentReplicaId : fallbackLabel;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -362,11 +362,52 @@ public static PubSubPositionGrpcWireFormat getPubSubPositionGrpcWireFormat(PubSu
* @param offset the fallback offset to use if deserialization fails or buffer is empty
* @param pubSubPositionDeserializer the deserializer to convert wire format to position
* @return a valid PubSubPosition, either deserialized or offset-based
* @deprecated Prefer {@link #deserializePositionWithOffsetFallback(ByteBuffer, long, PubSubPositionDeserializer, String)}
* so that a failure warning can be attributed to its source. This overload is retained only for tests and
* truly context-free callers; it logs {@code "N/A"} as the position source.
*/
@Deprecated
public static PubSubPosition deserializePositionWithOffsetFallback(
ByteBuffer wireFormatBytes,
long offset,
PubSubPositionDeserializer pubSubPositionDeserializer) {
return deserializePositionWithOffsetFallback(wireFormatBytes, offset, pubSubPositionDeserializer, null);
}

/**
* Same as {@link #deserializePositionWithOffsetFallback(ByteBuffer, long, PubSubPositionDeserializer)}, but also
* accepts a {@code positionSource} that is included in the warning logged on deserialization failure. This makes
* it possible to attribute "Failed to deserialize PubSubPosition" warnings to a specific origin instead of only
* the raw offset, which is otherwise insufficient to identify what emitted the malformed payload.
*
* <p>{@code positionSource} is intentionally a free-form string rather than a strongly-typed replica handle,
* because not every caller has a replica (store-version/partition) to point to. Callers should supply the most
* specific stable identifier they have, in this order of preference:
* <ol>
* <li>a canonical replica ID, e.g. {@code <store>_v<version>-<partition>} as produced by
* {@code Utils.getReplicaId}, when ingesting/consuming a specific partition</li>
* <li>a broker/field-qualified label (e.g. {@code "OffsetRecord.upstreamRealTimeTopicPubSubPosition[<broker>]"})
* when the caller only has a checkpoint field and no partition context</li>
* <li>a static, descriptive call-site label (e.g. {@code "OffsetRecord.lastProcessedVersionTopicPubSubPosition"})
* for utility code that has no per-call context at all</li>
* </ol>
* Fabricating an identifier that doesn't reflect real context (e.g. reusing an unrelated ID) defeats the purpose
* and must be avoided; passing {@code null} is only acceptable when no caller-supplied context is possible, in
* which case {@code "N/A"} is logged.
*
* @param wireFormatBytes the serialized position bytes (can be null or empty)
* @param offset the fallback offset to use if deserialization fails or buffer is empty
* @param pubSubPositionDeserializer the deserializer to convert wire format to position
* @param positionSource a stable, non-fabricated diagnostic label identifying the origin of {@code wireFormatBytes}
* for logging context; may be {@code null} when no context is available, in which case
* "N/A" is logged
* @return a valid PubSubPosition, either deserialized or offset-based
*/
public static PubSubPosition deserializePositionWithOffsetFallback(
ByteBuffer wireFormatBytes,
long offset,
PubSubPositionDeserializer pubSubPositionDeserializer,
String positionSource) {
// Fast path: nothing to deserialize
if (wireFormatBytes == null || !wireFormatBytes.hasRemaining()) {
return fromKafkaOffset(offset);
Expand All @@ -387,7 +428,8 @@ public static PubSubPosition deserializePositionWithOffsetFallback(
offset);
} catch (RuntimeException e) {
LOGGER.warn(
"Failed to deserialize PubSubPosition. Using offset-based position (offset={}, bufferRem={}, bufferCap={}).",
"Failed to deserialize PubSubPosition for: {}. Using offset-based position (offset={}, bufferRem={}, bufferCap={}).",
positionSource == null ? "N/A" : positionSource,
offset,
wireFormatBytes.remaining(),
wireFormatBytes.capacity(),
Expand Down
Loading
Loading