From d5af38c2429b60b77f54e0e2d0b170514c85dd1a Mon Sep 17 00:00:00 2001 From: artem Date: Tue, 10 Mar 2026 11:51:09 -0700 Subject: [PATCH 1/2] Add Memq Metadata Client Use one memq consumer Fix group_id make memqConsumer protected omit memq numbytes in metric Close iterator after persisting to messages list Add MemqConsumer PscSourceReaderMetrics support Update metrics capture Fix memq offset conversion Try to fix unit test failure Clean up MetricRegistry scheduler + change histograms to ExponentiallyDecayingReservoir Try to reduce hot path CPU util by caching offsets to prevent hashcode computation + hashmap lookup Bump memq.version to 1.0.2 --- psc-examples/pom.xml | 2 +- .../metrics/MemqSourceReaderMetricsUtil.java | 27 ++ .../metrics/PscSourceReaderMetrics.java | 59 ++++- .../PscTopicUriPartitionSplitReader.java | 5 +- psc-integration-test/pom.xml | 2 +- psc/pom.xml | 2 +- .../psc/common/TopicUriPartition.java | 14 +- ...taClientToMemqConsumerConfigConverter.java | 26 ++ .../consumer/PscConsumerMessagesIterable.java | 9 + .../psc/consumer/memq/MemqTopicUri.java | 2 +- .../psc/consumer/memq/PscMemqConsumer.java | 71 +++++- .../client/memq/PscMemqMetadataClient.java | 239 ++++++++++++++++++ .../PscMemqMetadataClientCreator.java | 37 +++ .../PscMetadataClientCreatorManager.java | 4 +- .../psc/metrics/PscMetricRegistryManager.java | 4 +- .../memq/TestMemqOffsetRoundTrip.java | 83 ++++++ 16 files changed, 559 insertions(+), 27 deletions(-) create mode 100644 psc-flink/src/main/java/com/pinterest/flink/connector/psc/source/metrics/MemqSourceReaderMetricsUtil.java create mode 100644 psc/src/main/java/com/pinterest/psc/config/PscMetadataClientToMemqConsumerConfigConverter.java create mode 100644 psc/src/main/java/com/pinterest/psc/metadata/client/memq/PscMemqMetadataClient.java create mode 100644 psc/src/main/java/com/pinterest/psc/metadata/creation/PscMemqMetadataClientCreator.java create mode 100644 psc/src/test/java/com/pinterest/psc/consumer/memq/TestMemqOffsetRoundTrip.java diff --git a/psc-examples/pom.xml b/psc-examples/pom.xml index 13215a76..7baa64e1 100644 --- a/psc-examples/pom.xml +++ b/psc-examples/pom.xml @@ -15,7 +15,7 @@ psc-examples - 0.2.21 + 1.0.2 diff --git a/psc-flink/src/main/java/com/pinterest/flink/connector/psc/source/metrics/MemqSourceReaderMetricsUtil.java b/psc-flink/src/main/java/com/pinterest/flink/connector/psc/source/metrics/MemqSourceReaderMetricsUtil.java new file mode 100644 index 00000000..0f9aa087 --- /dev/null +++ b/psc-flink/src/main/java/com/pinterest/flink/connector/psc/source/metrics/MemqSourceReaderMetricsUtil.java @@ -0,0 +1,27 @@ +package com.pinterest.flink.connector.psc.source.metrics; + +import com.pinterest.psc.common.TopicUriPartition; +import com.pinterest.psc.metrics.Metric; +import com.pinterest.psc.metrics.MetricName; + +import java.util.Map; +import java.util.function.Predicate; + +class MemqSourceReaderMetricsUtil { + + public static final String MEMQ_CONSUMER_METRIC_GROUP = "memq-consumer-metrics"; + public static final String BYTES_CONSUMED_TOTAL = "bytes.consumed.total"; + public static final String NOTIFICATION_RECORDS_LAG_MAX = "notification.records.lag.max"; + + protected static Predicate> createBytesConsumedFilter() { + return entry -> + entry.getKey().group().equals(MEMQ_CONSUMER_METRIC_GROUP) + && entry.getKey().name().equals(BYTES_CONSUMED_TOTAL); + } + + protected static Predicate> createRecordsLagFilter(TopicUriPartition tp) { + return entry -> + entry.getKey().group().equals(MEMQ_CONSUMER_METRIC_GROUP) + && entry.getKey().name().equals(NOTIFICATION_RECORDS_LAG_MAX); + } +} diff --git a/psc-flink/src/main/java/com/pinterest/flink/connector/psc/source/metrics/PscSourceReaderMetrics.java b/psc-flink/src/main/java/com/pinterest/flink/connector/psc/source/metrics/PscSourceReaderMetrics.java index e821b67d..a57d0686 100644 --- a/psc-flink/src/main/java/com/pinterest/flink/connector/psc/source/metrics/PscSourceReaderMetrics.java +++ b/psc-flink/src/main/java/com/pinterest/flink/connector/psc/source/metrics/PscSourceReaderMetrics.java @@ -25,6 +25,7 @@ import com.pinterest.psc.exception.ClientException; import com.pinterest.psc.metrics.Metric; import com.pinterest.psc.metrics.MetricName; +import java.util.Iterator; import org.apache.flink.annotation.PublicEvolving; import org.apache.flink.metrics.Counter; import org.apache.flink.metrics.MetricGroup; @@ -153,6 +154,19 @@ public void recordCurrentOffset(TopicUriPartition tp, long offset) { offsets.get(tp).currentOffset = offset; } + /** + * Returns the {@link Offset} tracker for the given partition, allowing callers to + * cache it and update offsets directly without repeated HashMap lookups. + * + * @param tp the topic partition to get the tracker for + * @return the Offset tracker + * @throws IllegalArgumentException if the partition is not tracked + */ + public Offset getOffsetTracker(TopicUriPartition tp) { + checkTopicPartitionTracked(tp); + return offsets.get(tp); + } + /** * Update the latest committed offset of the given {@link TopicUriPartition}. * @@ -180,8 +194,21 @@ public void recordFailedCommit() { * @param consumer Kafka consumer */ public void registerNumBytesIn(PscConsumer consumer) throws ClientException { - Predicate> filter = - KafkaSourceReaderMetricsUtil.createBytesConsumedFilter(); + String backendType = getBackendFromTags(consumer.metrics()); + Predicate> filter; + switch (backendType) { + case PscUtils.BACKEND_TYPE_KAFKA: + filter = KafkaSourceReaderMetricsUtil.createBytesConsumedFilter(); + break; + case PscUtils.BACKEND_TYPE_MEMQ: + filter = MemqSourceReaderMetricsUtil.createBytesConsumedFilter(); + break; + default: + LOG.warn( + "Unsupported backend type: \"{}\". Metric \"{}\" may not be reported correctly.", + backendType, MetricNames.IO_NUM_BYTES_IN); + return; + } this.bytesConsumedTotalMetric = MetricUtil.getPscMetric(consumer.metrics(), filter); } @@ -288,25 +315,35 @@ private void checkTopicPartitionTracked(TopicUriPartition tp) { case PscUtils.BACKEND_TYPE_KAFKA: filter = KafkaSourceReaderMetricsUtil.createRecordLagFilter(tp); break; + case PscUtils.BACKEND_TYPE_MEMQ: + filter = MemqSourceReaderMetricsUtil.createRecordsLagFilter(tp); + break; default: LOG.warn( - String.format( - "Unsupported backend type \"%s\". " - + "Metric \"%s\" may not be reported correctly. ", - backendType, MetricNames.PENDING_RECORDS)); + "Unsupported backend type \"{}\". Metric \"{}\" may not be reported correctly.", + backendType, MetricNames.PENDING_RECORDS); return null; } - return MetricUtil.getPscMetric(metrics, filter); + try { + return MetricUtil.getPscMetric(metrics, filter); + } catch (IllegalStateException e) { + LOG.debug("Metric not yet available for backend \"{}\", will retry on next poll cycle.", backendType); + return null; + } } private static String getBackendFromTags(Map metrics) { // sample the first entry to get the backend type - return metrics.keySet().iterator().next().tags().get("backend"); + Iterator it = metrics.keySet().iterator(); + if (!it.hasNext()) { + return "unknown"; + } + return it.next().tags().get("backend"); } - private static class Offset { - long currentOffset; - long committedOffset; + public static class Offset { + public long currentOffset; + public long committedOffset; Offset(long currentOffset, long committedOffset) { this.currentOffset = currentOffset; diff --git a/psc-flink/src/main/java/com/pinterest/flink/connector/psc/source/reader/PscTopicUriPartitionSplitReader.java b/psc-flink/src/main/java/com/pinterest/flink/connector/psc/source/reader/PscTopicUriPartitionSplitReader.java index 4e8adc54..9428f288 100644 --- a/psc-flink/src/main/java/com/pinterest/flink/connector/psc/source/reader/PscTopicUriPartitionSplitReader.java +++ b/psc-flink/src/main/java/com/pinterest/flink/connector/psc/source/reader/PscTopicUriPartitionSplitReader.java @@ -567,6 +567,7 @@ private static class PscPartitionSplitRecords private Iterator> recordIterator; private TopicUriPartition currentTopicPartition; private Long currentSplitStoppingOffset; + private PscSourceReaderMetrics.Offset currentOffsetTracker; private PscPartitionSplitRecords( PscConsumerMessagesIterable consumerMessagesIterable, PscSourceReaderMetrics metrics) { @@ -592,11 +593,13 @@ public String nextSplit() { recordIterator = consumerMessagesIterable.getMessagesForTopicUriPartition(currentTopicPartition).iterator(); currentSplitStoppingOffset = stoppingOffsets.getOrDefault(currentTopicPartition, Long.MAX_VALUE); + currentOffsetTracker = metrics.getOffsetTracker(currentTopicPartition); return currentTopicPartition.toString(); } else { currentTopicPartition = null; recordIterator = null; currentSplitStoppingOffset = null; + currentOffsetTracker = null; return null; } } @@ -612,7 +615,7 @@ public PscConsumerMessage nextRecordFromSplit() { final PscConsumerMessage message = recordIterator.next(); // Only emit records before stopping offset if (message.getMessageId().getOffset() < currentSplitStoppingOffset) { - metrics.recordCurrentOffset(currentTopicPartition, message.getMessageId().getOffset()); + currentOffsetTracker.currentOffset = message.getMessageId().getOffset(); return message; } } diff --git a/psc-integration-test/pom.xml b/psc-integration-test/pom.xml index 22e0d292..ed51b94c 100644 --- a/psc-integration-test/pom.xml +++ b/psc-integration-test/pom.xml @@ -14,7 +14,7 @@ 3.4.0 - 0.2.21 + 1.0.2 diff --git a/psc/pom.xml b/psc/pom.xml index f4a65bf2..cc06c6a6 100644 --- a/psc/pom.xml +++ b/psc/pom.xml @@ -15,7 +15,7 @@ 3.4.0 - 0.2.21 + 1.0.2 diff --git a/psc/src/main/java/com/pinterest/psc/common/TopicUriPartition.java b/psc/src/main/java/com/pinterest/psc/common/TopicUriPartition.java index 39b8a1c0..333b56cc 100644 --- a/psc/src/main/java/com/pinterest/psc/common/TopicUriPartition.java +++ b/psc/src/main/java/com/pinterest/psc/common/TopicUriPartition.java @@ -12,6 +12,7 @@ public class TopicUriPartition implements Comparable, Seriali private final String topicUriStr; private final int partition; private TopicUri backendTopicUri; + private transient int cachedHashCode; /** * Builds a TopicUriPartition instance with the default partition value (-1). This is meant to be used in @@ -53,6 +54,7 @@ public TopicUriPartition(TopicUri topicUri, int partition) { protected void setTopicUri(TopicUri backendTopicUri) { this.backendTopicUri = backendTopicUri; + this.cachedHashCode = 0; } /** @@ -106,10 +108,14 @@ public boolean equals(Object other) { @Override public int hashCode() { - int result = topicUriStr.hashCode(); - result = 31 * result + (backendTopicUri == null ? 0 : backendTopicUri.hashCode()); - result = 31 * result + partition; - return result; + int h = cachedHashCode; + if (h == 0) { + h = topicUriStr.hashCode(); + h = 31 * h + (backendTopicUri == null ? 0 : backendTopicUri.hashCode()); + h = 31 * h + partition; + cachedHashCode = h; + } + return h; } @Override diff --git a/psc/src/main/java/com/pinterest/psc/config/PscMetadataClientToMemqConsumerConfigConverter.java b/psc/src/main/java/com/pinterest/psc/config/PscMetadataClientToMemqConsumerConfigConverter.java new file mode 100644 index 00000000..2c5ce98e --- /dev/null +++ b/psc/src/main/java/com/pinterest/psc/config/PscMetadataClientToMemqConsumerConfigConverter.java @@ -0,0 +1,26 @@ +package com.pinterest.psc.config; + +import com.pinterest.memq.client.commons.ConsumerConfigs; +import com.pinterest.psc.common.TopicUri; + +import java.util.HashMap; +import java.util.Map; +import java.util.Properties; + +public class PscMetadataClientToMemqConsumerConfigConverter extends PscMetadataClientToBackendMetatadataClientConfigCoverter { + @Override + protected Map getConfigConverterMap() { + return new HashMap() { + private static final long serialVersionUID = 1L; + + { + put(PscConfiguration.PSC_METADATA_CLIENT_ID, ConsumerConfigs.CLIENT_ID); + } + }; + } + + @Override + public Properties convert(PscConfigurationInternal pscConfigurationInternal, TopicUri topicUri) { + return super.convert(pscConfigurationInternal, topicUri); + } +} diff --git a/psc/src/main/java/com/pinterest/psc/consumer/PscConsumerMessagesIterable.java b/psc/src/main/java/com/pinterest/psc/consumer/PscConsumerMessagesIterable.java index 211a7511..9f1ace39 100644 --- a/psc/src/main/java/com/pinterest/psc/consumer/PscConsumerMessagesIterable.java +++ b/psc/src/main/java/com/pinterest/psc/consumer/PscConsumerMessagesIterable.java @@ -1,7 +1,9 @@ package com.pinterest.psc.consumer; import com.pinterest.psc.common.TopicUriPartition; +import com.pinterest.psc.logging.PscLogger; +import java.io.IOException; import java.util.ArrayList; import java.util.HashMap; import java.util.Iterator; @@ -13,11 +15,18 @@ public class PscConsumerMessagesIterable implements Iterable> { + private static final PscLogger logger = PscLogger.getLogger(PscConsumerMessagesIterable.class); + List> messages; Map>> messagesByTopicUriPartition; public PscConsumerMessagesIterable(PscConsumerPollMessageIterator iterator) { this.messages = iterator.asList(); + try { + iterator.close(); + } catch (IOException e) { + logger.warn("Failed to close poll message iterator", e); + } this.messagesByTopicUriPartition = new HashMap<>(); for (PscConsumerMessage message : messages) { TopicUriPartition topicUriPartition = message.getMessageId().getTopicUriPartition(); diff --git a/psc/src/main/java/com/pinterest/psc/consumer/memq/MemqTopicUri.java b/psc/src/main/java/com/pinterest/psc/consumer/memq/MemqTopicUri.java index d25c1a65..3a1760a5 100644 --- a/psc/src/main/java/com/pinterest/psc/consumer/memq/MemqTopicUri.java +++ b/psc/src/main/java/com/pinterest/psc/consumer/memq/MemqTopicUri.java @@ -9,7 +9,7 @@ public class MemqTopicUri extends BaseTopicUri { public static final String PLAINTEXT_PROTOCOL = "plaintext"; public static final String SECURE_PROTOCOL = "secure"; - MemqTopicUri(TopicUri topicUri) { + public MemqTopicUri(TopicUri topicUri) { super(topicUri); } diff --git a/psc/src/main/java/com/pinterest/psc/consumer/memq/PscMemqConsumer.java b/psc/src/main/java/com/pinterest/psc/consumer/memq/PscMemqConsumer.java index e93dd667..33af7283 100644 --- a/psc/src/main/java/com/pinterest/psc/consumer/memq/PscMemqConsumer.java +++ b/psc/src/main/java/com/pinterest/psc/consumer/memq/PscMemqConsumer.java @@ -29,6 +29,7 @@ import com.pinterest.psc.exception.consumer.ConsumerException; import com.pinterest.psc.exception.consumer.WakeupException; import com.pinterest.psc.logging.PscLogger; +import com.pinterest.psc.common.PscUtils; import com.pinterest.psc.metrics.Metric; import com.pinterest.psc.metrics.MetricName; import com.pinterest.psc.metrics.PscMetricRegistryManager; @@ -55,6 +56,7 @@ public class PscMemqConsumer extends PscBackendConsumer { public static final String END_OF_BATCH_EVENT = "end_of_batch"; + private static final String MEMQ_CONSUMER_METRIC_GROUP = "memq-consumer-metrics"; private static final PscLogger logger = PscLogger.getLogger(PscMemqConsumer.class); @VisibleForTesting @@ -599,6 +601,7 @@ public void wakeup() { public void close() throws ConsumerException { if (memqConsumer == null) throw new ConsumerException("[Memq] Consumer is not initialized prior to call to close()."); + scheduler.shutdown(); currentSubscription.clear(); try { memqConsumer.close(); @@ -640,7 +643,8 @@ public Map startOffsets(Set topicUri Map startOffsets = memqConsumer .getEarliestOffsets(partitionToTopicUriPartition.keySet()); return startOffsets.entrySet().stream().collect(Collectors - .toMap(entry -> partitionToTopicUriPartition.get(entry.getKey()), Map.Entry::getValue)); + .toMap(entry -> partitionToTopicUriPartition.get(entry.getKey()), + entry -> kafkaOffsetToComposite(entry.getValue()))); } @Override @@ -654,7 +658,8 @@ public Map endOffsets(Set topicUriPa Map endOffsets = memqConsumer .getLatestOffsets(partitionToTopicUriPartition.keySet()); return endOffsets.entrySet().stream().collect(Collectors - .toMap(entry -> partitionToTopicUriPartition.get(entry.getKey()), Map.Entry::getValue)); + .toMap(entry -> partitionToTopicUriPartition.get(entry.getKey()), + entry -> kafkaOffsetToComposite(entry.getValue()))); } @Override @@ -714,7 +719,57 @@ public PscConfiguration getConfiguration() { @Override public Map metrics() throws ConsumerException { - return Collections.emptyMap(); + if (memqConsumer == null) { + return Collections.emptyMap(); + } + + MetricRegistry registry = memqConsumer.getMetricRegistry(); + if (registry == null) { + return Collections.emptyMap(); + } + + Map result = new HashMap<>(); + for (Map.Entry entry : registry.getMetrics().entrySet()) { + String name = entry.getKey(); + com.codahale.metrics.Metric dropwizardMetric = entry.getValue(); + + Map tags = new HashMap<>(); + tags.put("backend", PscUtils.BACKEND_TYPE_MEMQ); + + MetricName metricName = new MetricName(name, MEMQ_CONSUMER_METRIC_GROUP, "", tags); + result.put(metricName, new LiveDropwizardMetric(metricName, dropwizardMetric)); + } + + return result; + } + + /** + * A PSC Metric backed by a live Dropwizard metric reference. + * Each call to {@link #metricValue()} reads the current value from the + * underlying Dropwizard metric rather than returning a stale snapshot. + */ + private static class LiveDropwizardMetric extends Metric { + private final com.codahale.metrics.Metric dropwizardMetric; + + LiveDropwizardMetric(MetricName metricName, com.codahale.metrics.Metric dropwizardMetric) { + super(metricName, null); + this.dropwizardMetric = dropwizardMetric; + } + + @Override + public Object metricValue() { + if (dropwizardMetric instanceof Counter) + return ((Counter) dropwizardMetric).getCount(); + if (dropwizardMetric instanceof Gauge) + return ((Gauge) dropwizardMetric).getValue(); + if (dropwizardMetric instanceof Meter) + return ((Meter) dropwizardMetric).getCount(); + if (dropwizardMetric instanceof Histogram) + return ((Histogram) dropwizardMetric).getSnapshot().getMax(); + if (dropwizardMetric instanceof Timer) + return ((Timer) dropwizardMetric).getSnapshot().getMax(); + return -1L; + } } /** @@ -746,6 +801,16 @@ private boolean isCurrentTopicPartition(TopicUriPartition topicUriPartition) { return this.currentSubscription.contains(topicUriPartition.getTopicUri()) || this.currentAssignment.contains(topicUriPartition); } + /** + * Converts a raw Kafka notification offset to a composite MemqOffset (with message offset 0). + * All offsets exposed by PscMemqConsumer must be in composite format so that + * {@link #seekToOffset} can correctly decode them back via + * {@link MemqOffset#convertPscOffsetToMemqOffset}. + */ + private static long kafkaOffsetToComposite(long kafkaOffset) { + return new MemqOffset(kafkaOffset, 0).toLong(); + } + private MemqConsumer getMetadataConsumer(TopicUri topicUri) throws ConsumerException { try { Properties tmpProps = new Properties(properties); diff --git a/psc/src/main/java/com/pinterest/psc/metadata/client/memq/PscMemqMetadataClient.java b/psc/src/main/java/com/pinterest/psc/metadata/client/memq/PscMemqMetadataClient.java new file mode 100644 index 00000000..f26bf516 --- /dev/null +++ b/psc/src/main/java/com/pinterest/psc/metadata/client/memq/PscMemqMetadataClient.java @@ -0,0 +1,239 @@ +package com.pinterest.psc.metadata.client.memq; + +import com.pinterest.memq.client.commons.ConsumerConfigs; +import com.pinterest.memq.client.commons.serde.ByteArrayDeserializer; +import com.pinterest.memq.client.consumer.MemqConsumer; +import com.pinterest.psc.common.BaseTopicUri; +import com.pinterest.psc.common.TopicRn; +import com.pinterest.psc.common.TopicUri; +import com.pinterest.psc.common.TopicUriPartition; +import com.pinterest.psc.config.PscConfigurationInternal; +import com.pinterest.psc.config.PscMetadataClientToMemqConsumerConfigConverter; +import com.pinterest.psc.consumer.memq.MemqOffset; +import com.pinterest.psc.consumer.memq.MemqTopicUri; +import com.pinterest.psc.environment.Environment; +import com.pinterest.psc.exception.startup.ConfigurationException; +import com.pinterest.psc.logging.PscLogger; +import com.pinterest.psc.metadata.MetadataUtils; +import com.pinterest.psc.metadata.TopicUriMetadata; +import com.pinterest.psc.metadata.client.PscBackendMetadataClient; +import com.pinterest.psc.metadata.client.PscMetadataClient; + +import java.io.IOException; +import java.time.Duration; +import java.util.ArrayList; +import java.util.Collection; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Properties; +import java.util.Set; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.TimeoutException; + +/** + * A Memq-specific implementation of the {@link PscBackendMetadataClient}. + * Uses a {@link MemqConsumer} to query metadata since Memq does not have a dedicated admin client. + */ +public class PscMemqMetadataClient extends PscBackendMetadataClient { + + private static final PscLogger logger = PscLogger.getLogger(PscMemqMetadataClient.class); + protected MemqConsumer memqConsumer; + + @Override + public void initialize( + TopicUri topicUri, + Environment env, + PscConfigurationInternal pscConfigurationInternal + ) throws ConfigurationException { + super.initialize(topicUri, env, pscConfigurationInternal); + Properties properties = new PscMetadataClientToMemqConsumerConfigConverter() + .convert(pscConfigurationInternal, topicUri); + properties.setProperty(ConsumerConfigs.BOOTSTRAP_SERVERS, discoveryConfig.getConnect()); + properties.setProperty(ConsumerConfigs.CLIENT_ID, + pscConfigurationInternal.getMetadataClientId()); + properties.setProperty(ConsumerConfigs.GROUP_ID, + pscConfigurationInternal.getMetadataClientId()); + properties.setProperty(ConsumerConfigs.KEY_DESERIALIZER_CLASS_KEY, + ByteArrayDeserializer.class.getName()); + properties.put(ConsumerConfigs.KEY_DESERIALIZER_CLASS_CONFIGS_KEY, new Properties()); + properties.setProperty(ConsumerConfigs.VALUE_DESERIALIZER_CLASS_KEY, + ByteArrayDeserializer.class.getName()); + properties.put(ConsumerConfigs.VALUE_DESERIALIZER_CLASS_CONFIGS_KEY, new Properties()); + properties.setProperty(ConsumerConfigs.DIRECT_CONSUMER, "false"); + try { + this.memqConsumer = new MemqConsumer<>(properties); + } catch (Exception e) { + throw new ConfigurationException("Failed to create Memq consumer for metadata client", e); + } + logger.info("Initialized PscMemqMetadataClient with properties: " + properties); + } + + @Override + public List listTopicRns(Duration duration) + throws ExecutionException, InterruptedException, TimeoutException { + throw new UnsupportedOperationException( + "[Memq] Listing all topics is not supported by the Memq backend."); + } + + @Override + public Map describeTopicUris( + Collection topicUris, + Duration duration + ) throws ExecutionException, InterruptedException, TimeoutException { + Map result = new HashMap<>(); + for (TopicUri tu : topicUris) { + subscribe(tu.getTopic()); + List partitions = memqConsumer.getPartition(); + List topicUriPartitions = new ArrayList<>(); + for (int partition : partitions) { + topicUriPartitions.add(createMemqTopicUriPartition(tu, partition)); + } + result.put(tu, new TopicUriMetadata(tu, topicUriPartitions)); + } + return result; + } + + @Override + public Map listOffsets( + Map topicUriPartitionsAndOptions, + Duration duration + ) throws ExecutionException, InterruptedException, TimeoutException { + Map> earliestByTopic = new HashMap<>(); + Map> latestByTopic = new HashMap<>(); + + for (Map.Entry entry : + topicUriPartitionsAndOptions.entrySet()) { + TopicUriPartition tup = entry.getKey(); + String topic = tup.getTopicUri().getTopic(); + + if (entry.getValue() == PscMetadataClient.MetadataClientOption.OFFSET_SPEC_EARLIEST) { + earliestByTopic.computeIfAbsent(topic, k -> new HashSet<>()).add(tup.getPartition()); + } else if (entry.getValue() == PscMetadataClient.MetadataClientOption.OFFSET_SPEC_LATEST) { + latestByTopic.computeIfAbsent(topic, k -> new HashSet<>()).add(tup.getPartition()); + } else { + throw new IllegalArgumentException( + "Unsupported MetadataClientOption for listOffsets(): " + entry.getValue()); + } + } + + Map result = new HashMap<>(); + Set allTopics = new HashSet<>(); + allTopics.addAll(earliestByTopic.keySet()); + allTopics.addAll(latestByTopic.keySet()); + + for (String topic : allTopics) { + subscribe(topic); + + Set earliestPartitions = earliestByTopic.getOrDefault(topic, new HashSet<>()); + if (!earliestPartitions.isEmpty()) { + Map offsets = memqConsumer.getEarliestOffsets(earliestPartitions); + for (Map.Entry e : offsets.entrySet()) { + TopicRn topicRn = MetadataUtils.createTopicRn(topicUri, topic); + result.put(createMemqTopicUriPartition(topicRn, e.getKey()), kafkaOffsetToComposite(e.getValue())); + } + } + + Set latestPartitions = latestByTopic.getOrDefault(topic, new HashSet<>()); + if (!latestPartitions.isEmpty()) { + Map offsets = memqConsumer.getLatestOffsets(latestPartitions); + for (Map.Entry e : offsets.entrySet()) { + TopicRn topicRn = MetadataUtils.createTopicRn(topicUri, topic); + result.put(createMemqTopicUriPartition(topicRn, e.getKey()), kafkaOffsetToComposite(e.getValue())); + } + } + } + return result; + } + + @Override + public Map listOffsetsForTimestamps( + Map topicUriPartitionsAndTimes, + Duration duration + ) throws ExecutionException, InterruptedException, TimeoutException { + Map> timestampsByTopic = new HashMap<>(); + + for (Map.Entry entry : topicUriPartitionsAndTimes.entrySet()) { + TopicUriPartition tup = entry.getKey(); + String topic = tup.getTopicUri().getTopic(); + timestampsByTopic.computeIfAbsent(topic, k -> new HashMap<>()) + .put(tup.getPartition(), entry.getValue()); + } + + Map result = new HashMap<>(); + for (Map.Entry> entry : timestampsByTopic.entrySet()) { + String topic = entry.getKey(); + subscribe(topic); + Map offsets = memqConsumer.offsetsOfTimestamps(entry.getValue()); + for (Map.Entry offsetEntry : offsets.entrySet()) { + TopicRn topicRn = MetadataUtils.createTopicRn(topicUri, topic); + result.put( + createMemqTopicUriPartition(topicRn, offsetEntry.getKey()), + kafkaOffsetToComposite(offsetEntry.getValue()) + ); + } + } + return result; + } + + @Override + public Map listOffsetsForConsumerGroup( + String consumerGroupId, + Collection topicUriPartitions, + Duration duration + ) throws ExecutionException, InterruptedException, TimeoutException { + Map> partitionsByTopic = new HashMap<>(); + for (TopicUriPartition tup : topicUriPartitions) { + String topic = tup.getTopicUri().getTopic(); + partitionsByTopic.computeIfAbsent(topic, k -> new HashSet<>()).add(tup.getPartition()); + } + + Map result = new HashMap<>(); + for (Map.Entry> entry : partitionsByTopic.entrySet()) { + String topic = entry.getKey(); + subscribe(topic); + for (int partition : entry.getValue()) { + long committedOffset = memqConsumer.committed(partition); + if (committedOffset == -1L) { + logger.warn( + "Consumer group {} has no committed offset for topic {} partition {}", + consumerGroupId, topic, partition + ); + continue; + } + TopicRn topicRn = MetadataUtils.createTopicRn(topicUri, topic); + result.put(createMemqTopicUriPartition(topicRn, partition), kafkaOffsetToComposite(committedOffset)); + } + } + return result; + } + + @Override + public void close() throws IOException { + if (memqConsumer != null) + memqConsumer.close(); + logger.info("Closed PscMemqMetadataClient"); + } + + private void subscribe(String topic) throws ExecutionException { + try { + memqConsumer.subscribe(topic); + } catch (Exception e) { + throw new ExecutionException("Failed to subscribe to Memq topic " + topic, e); + } + } + + private TopicUriPartition createMemqTopicUriPartition(TopicRn topicRn, int partition) { + return new TopicUriPartition( + new MemqTopicUri(new BaseTopicUri(topicUri.getProtocol(), topicRn)), partition); + } + + private TopicUriPartition createMemqTopicUriPartition(TopicUri topicUri, int partition) { + return new TopicUriPartition(new MemqTopicUri(topicUri), partition); + } + + private static long kafkaOffsetToComposite(long kafkaOffset) { + return new MemqOffset(kafkaOffset, 0).toLong(); + } +} diff --git a/psc/src/main/java/com/pinterest/psc/metadata/creation/PscMemqMetadataClientCreator.java b/psc/src/main/java/com/pinterest/psc/metadata/creation/PscMemqMetadataClientCreator.java new file mode 100644 index 00000000..d5645e1b --- /dev/null +++ b/psc/src/main/java/com/pinterest/psc/metadata/creation/PscMemqMetadataClientCreator.java @@ -0,0 +1,37 @@ +package com.pinterest.psc.metadata.creation; + +import com.pinterest.psc.common.PscUtils; +import com.pinterest.psc.common.TopicUri; +import com.pinterest.psc.config.PscConfigurationInternal; +import com.pinterest.psc.consumer.memq.MemqTopicUri; +import com.pinterest.psc.environment.Environment; +import com.pinterest.psc.exception.startup.ConfigurationException; +import com.pinterest.psc.exception.startup.TopicUriSyntaxException; +import com.pinterest.psc.logging.PscLogger; +import com.pinterest.psc.metadata.client.memq.PscMemqMetadataClient; + +/** + * A class that creates a {@link com.pinterest.psc.metadata.client.PscBackendMetadataClient} for Memq. + */ +@PscMetadataClientCreatorPlugin(backend = PscUtils.BACKEND_TYPE_MEMQ, priority = 1) +public class PscMemqMetadataClientCreator extends PscBackendMetadataClientCreator { + + private static final PscLogger logger = PscLogger.getLogger(PscMemqMetadataClientCreator.class); + + @Override + public PscMemqMetadataClient create(Environment env, PscConfigurationInternal pscConfigurationInternal, TopicUri clusterUri) throws ConfigurationException { + logger.info("Creating Memq metadata client for clusterUri: " + clusterUri); + PscMemqMetadataClient pscMemqMetadataClient = new PscMemqMetadataClient(); + pscMemqMetadataClient.initialize( + clusterUri, + env, + pscConfigurationInternal + ); + return pscMemqMetadataClient; + } + + @Override + public TopicUri validateBackendTopicUri(TopicUri topicUri) throws TopicUriSyntaxException { + return MemqTopicUri.validate(topicUri); + } +} diff --git a/psc/src/main/java/com/pinterest/psc/metadata/creation/PscMetadataClientCreatorManager.java b/psc/src/main/java/com/pinterest/psc/metadata/creation/PscMetadataClientCreatorManager.java index 0834411e..5884f29f 100644 --- a/psc/src/main/java/com/pinterest/psc/metadata/creation/PscMetadataClientCreatorManager.java +++ b/psc/src/main/java/com/pinterest/psc/metadata/creation/PscMetadataClientCreatorManager.java @@ -7,10 +7,10 @@ import java.util.ArrayList; import java.util.Collections; import java.util.Comparator; -import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Set; +import java.util.TreeMap; /** * Manages the different {@link PscBackendMetadataClientCreator} implementations and provides a registry of them. @@ -28,7 +28,7 @@ public class PscMetadataClientCreatorManager { private static Map> findAndRegisterMetadataClientCreators(String packageName) { synchronized (PscUtils.lock) { - Map> backendCreatorRegistry = new HashMap<>(); + Map> backendCreatorRegistry = new TreeMap<>(); Reflections reflections = new Reflections(packageName.trim()); Set> annotatedClasses = reflections.getTypesAnnotatedWith(PscMetadataClientCreatorPlugin.class); for (Class annotatedClass : annotatedClasses) { diff --git a/psc/src/main/java/com/pinterest/psc/metrics/PscMetricRegistryManager.java b/psc/src/main/java/com/pinterest/psc/metrics/PscMetricRegistryManager.java index 7598701c..a66bdc67 100644 --- a/psc/src/main/java/com/pinterest/psc/metrics/PscMetricRegistryManager.java +++ b/psc/src/main/java/com/pinterest/psc/metrics/PscMetricRegistryManager.java @@ -6,7 +6,7 @@ import com.codahale.metrics.MetricFilter; import com.codahale.metrics.MetricRegistry; import com.codahale.metrics.ScheduledReporter; -import com.codahale.metrics.SlidingTimeWindowArrayReservoir; +import com.codahale.metrics.ExponentiallyDecayingReservoir; import com.codahale.metrics.Snapshot; import com.codahale.metrics.jvm.CachedThreadStatesGaugeSet; import com.codahale.metrics.jvm.GarbageCollectorMetricSet; @@ -224,7 +224,7 @@ public void updateHistogramMetric(TopicUri topicUri, if (metricRegistry != null) { metricRegistry.histogram(metricKey, () -> new Histogram( - new SlidingTimeWindowArrayReservoir(1, TimeUnit.MINUTES) + new ExponentiallyDecayingReservoir() ) ).update(metricValue); } diff --git a/psc/src/test/java/com/pinterest/psc/consumer/memq/TestMemqOffsetRoundTrip.java b/psc/src/test/java/com/pinterest/psc/consumer/memq/TestMemqOffsetRoundTrip.java new file mode 100644 index 00000000..277a6f7c --- /dev/null +++ b/psc/src/test/java/com/pinterest/psc/consumer/memq/TestMemqOffsetRoundTrip.java @@ -0,0 +1,83 @@ +package com.pinterest.psc.consumer.memq; + +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +/** + * Verifies that raw Kafka notification offsets survive a round-trip through + * the composite MemqOffset encoding used by PscMemqConsumer. + * + * Bug context: offsetsForTimes / endOffsets / startOffsets were returning raw + * Kafka offsets, but seekToOffset decodes them as composite MemqOffsets + * (bit-shifting right by 19). Without wrapping via kafkaOffsetToComposite, + * a raw offset like 2205646 would be decoded as batch=4 instead of batch=2205646. + */ +public class TestMemqOffsetRoundTrip { + + @Test + public void testRawKafkaOffsetRoundTrips() { + long rawKafkaOffset = 2205646L; + + long composite = new MemqOffset(rawKafkaOffset, 0).toLong(); + MemqOffset decoded = MemqOffset.convertPscOffsetToMemqOffset(composite); + + assertEquals(rawKafkaOffset, decoded.getBatchOffset()); + assertEquals(0, decoded.getMessageOffset()); + } + + @Test + public void testRawOffsetWithoutEncodingIsCorrupted() { + long rawKafkaOffset = 2205646L; + + // Decoding a raw offset directly (the old bug) produces wrong batch offset + MemqOffset decoded = MemqOffset.convertPscOffsetToMemqOffset(rawKafkaOffset); + + // 2205646 >>> 19 = 4, not 2205646 + assertEquals(4, decoded.getBatchOffset(), + "Raw offset decoded without encoding should lose upper bits"); + } + + @Test + public void testSmallOffsetRoundTrips() { + long rawKafkaOffset = 42L; + + long composite = new MemqOffset(rawKafkaOffset, 0).toLong(); + MemqOffset decoded = MemqOffset.convertPscOffsetToMemqOffset(composite); + + assertEquals(rawKafkaOffset, decoded.getBatchOffset()); + assertEquals(0, decoded.getMessageOffset()); + } + + @Test + public void testZeroOffsetRoundTrips() { + long composite = new MemqOffset(0, 0).toLong(); + MemqOffset decoded = MemqOffset.convertPscOffsetToMemqOffset(composite); + + assertEquals(0, decoded.getBatchOffset()); + assertEquals(0, decoded.getMessageOffset()); + } + + @Test + public void testLargeOffsetRoundTrips() { + long rawKafkaOffset = 10_000_000L; + + long composite = new MemqOffset(rawKafkaOffset, 0).toLong(); + MemqOffset decoded = MemqOffset.convertPscOffsetToMemqOffset(composite); + + assertEquals(rawKafkaOffset, decoded.getBatchOffset()); + assertEquals(0, decoded.getMessageOffset()); + } + + @Test + public void testCompositeWithMessageOffsetRoundTrips() { + long batchOffset = 500L; + int messageOffset = 1234; + + long composite = new MemqOffset(batchOffset, messageOffset).toLong(); + MemqOffset decoded = MemqOffset.convertPscOffsetToMemqOffset(composite); + + assertEquals(batchOffset, decoded.getBatchOffset()); + assertEquals(messageOffset, decoded.getMessageOffset()); + } +} From b89228ec0ce73af186f0ddc8af7d9982cb9af11d Mon Sep 17 00:00:00 2001 From: Jeff Xiang Date: Wed, 15 Apr 2026 10:43:34 -0700 Subject: [PATCH 2/2] Handle null backend types --- .../connector/psc/source/metrics/PscSourceReaderMetrics.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/psc-flink/src/main/java/com/pinterest/flink/connector/psc/source/metrics/PscSourceReaderMetrics.java b/psc-flink/src/main/java/com/pinterest/flink/connector/psc/source/metrics/PscSourceReaderMetrics.java index a57d0686..644bd262 100644 --- a/psc-flink/src/main/java/com/pinterest/flink/connector/psc/source/metrics/PscSourceReaderMetrics.java +++ b/psc-flink/src/main/java/com/pinterest/flink/connector/psc/source/metrics/PscSourceReaderMetrics.java @@ -333,12 +333,12 @@ private void checkTopicPartitionTracked(TopicUriPartition tp) { } private static String getBackendFromTags(Map metrics) { - // sample the first entry to get the backend type Iterator it = metrics.keySet().iterator(); if (!it.hasNext()) { return "unknown"; } - return it.next().tags().get("backend"); + String backend = it.next().tags().get("backend"); + return backend != null ? backend : "unknown"; } public static class Offset {