[
https://issues.apache.org/jira/browse/FLINK-4460?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=15925917#comment-15925917
]
ASF GitHub Bot commented on FLINK-4460:
---------------------------------------
Github user kl0u commented on a diff in the pull request:
https://github.com/apache/flink/pull/3484#discussion_r106138473
--- Diff:
flink-streaming-java/src/main/java/org/apache/flink/streaming/api/datastream/SingleOutputStreamOperator.java
---
@@ -416,4 +428,26 @@ private boolean canBeParallel() {
transformation.setSlotSharingGroup(slotSharingGroup);
return this;
}
+
+ /**
+ * Gets the {@link DataStream} that contains the elements that are
emitted from an operation
+ * into the side output with the given {@link OutputTag}.
+ *
+ * @see
org.apache.flink.streaming.api.functions.ProcessFunction.Context#output(OutputTag,
Object)
+ */
+ public <X> DataStream<X> getSideOutput(OutputTag<X> sideOutputTag){
+ sideOutputTag = clean(sideOutputTag);
+
+ TypeInformation<?> type =
requestedSideOutputs.get(sideOutputTag);
+ if (type != null && !type.equals(sideOutputTag.getTypeInfo())) {
+ throw new UnsupportedOperationException("A side output
with a matching id was " +
+ "already requested with a different
type. This is not allowed, side output " +
+ "ids need to be unique.");
+ }
+
+ requestedSideOutputs.put(sideOutputTag,
sideOutputTag.getTypeInfo());
+
--- End diff --
The `requireNotNull` should be in the beginning of the method.
> Side Outputs in Flink
> ---------------------
>
> Key: FLINK-4460
> URL: https://issues.apache.org/jira/browse/FLINK-4460
> Project: Flink
> Issue Type: New Feature
> Components: Core, DataStream API
> Affects Versions: 1.2.0, 1.1.3
> Reporter: Chen Qin
> Assignee: Chen Qin
> Labels: latearrivingevents, sideoutput
>
> https://docs.google.com/document/d/1vg1gpR8JL4dM07Yu4NyerQhhVvBlde5qdqnuJv4LcV4/edit?usp=sharing
--
This message was sent by Atlassian JIRA
(v6.3.15#6346)