[ https://issues.apache.org/jira/browse/BEAM-6161?focusedWorklogId=182770&page=com.atlassian.jira.plugin.system.issuetabpanels:worklog-tabpanel#worklog-182770 ]
ASF GitHub Bot logged work on BEAM-6161: ---------------------------------------- Author: ASF GitHub Bot Created on: 09/Jan/19 00:18 Start Date: 09/Jan/19 00:18 Worklog Time Spent: 10m Work Description: Ardagan commented on pull request #7272: [BEAM-6161] Introduce PCollectionConsumerRegistry and add ElementCoun… URL: https://github.com/apache/beam/pull/7272#discussion_r246208702 ########## File path: sdks/java/harness/src/main/java/org/apache/beam/fn/harness/data/PCollectionConsumerRegistry.java ########## @@ -33,32 +36,65 @@ public class PCollectionConsumerRegistry { private ListMultimap<String, FnDataReceiver<WindowedValue<?>>> pCollectionIdsToConsumers; + private Map<String, ElementCountFnDataReceiver> pCollectionIdsToWrappedConsumer; public PCollectionConsumerRegistry() { pCollectionIdsToConsumers = ArrayListMultimap.create(); + pCollectionIdsToWrappedConsumer = new HashMap<String, ElementCountFnDataReceiver>(); } - public <T> FnDataReceiver<WindowedValue<T>> registerAndWrap( - String pCollectionId, FnDataReceiver<WindowedValue<T>> consumer) { - // TODO(ajamato): Test multiple consumers of the same pcollection before merging to master. - FnDataReceiver<WindowedValue<T>> wrappedReceiver = - new ElementCountFnDataReceiver<T>(consumer, pCollectionId); - pCollectionIdsToConsumers.put(pCollectionId, (FnDataReceiver) wrappedReceiver); - return wrappedReceiver; + public <T> void register(String pCollectionId, FnDataReceiver<WindowedValue<T>> consumer) { + // Just save these consumers for now, but package them up later with an + // ElementCountFnDataReceiver and possibly a MultiplexingFnDataReceiver + // if there are multiple consumers. + ElementCountFnDataReceiver wrappedConsumer = + pCollectionIdsToWrappedConsumer.getOrDefault(pCollectionId, null); + if (wrappedConsumer != null) { + throw new RuntimeException( Review comment: ack ---------------------------------------------------------------- This is an automated message from the Apache Git Service. To respond to the message, please log on GitHub and use the URL above to go to the specific comment. For queries about this service, please contact Infrastructure at: us...@infra.apache.org Issue Time Tracking ------------------- Worklog Id: (was: 182770) Time Spent: 6.5h (was: 6h 20m) > Add ElementCount MonitoringInfos for the Java SDK > ------------------------------------------------- > > Key: BEAM-6161 > URL: https://issues.apache.org/jira/browse/BEAM-6161 > Project: Beam > Issue Type: New Feature > Components: java-fn-execution, sdk-java-harness > Reporter: Alex Amato > Assignee: Alex Amato > Priority: Major > Time Spent: 6.5h > Remaining Estimate: 0h > -- This message was sent by Atlassian JIRA (v7.6.3#76005)