rkhachatryan commented on code in PR #28459:
URL: https://github.com/apache/flink/pull/28459#discussion_r3785535249
##########
flink-runtime/src/main/java/org/apache/flink/streaming/api/graph/AdaptiveGraphManager.java:
##########
@@ -447,6 +450,50 @@ private void
connectToFinishedUpStreamVertex(JobVertexBuildContext jobVertexBuil
}
}
+ private List<StreamEdge> getTransitiveInEdgesInOrder(
+ List<StreamEdge> transitiveInEdges, JobVertexBuildContext
jobVertexBuildContext) {
+ final List<StreamEdge> transitiveInEdgesInOrder = new
ArrayList<>(transitiveInEdges);
+ final List<StreamEdge> uidTransitiveInEdges =
+ transitiveInEdgesInOrder.stream()
+ .filter(this::hasUidBackedUpstream)
+ .sorted(
+ Comparator.comparing(
+ inEdge ->
+ getStartNodeJobVertexId(
+ inEdge,
jobVertexBuildContext)))
+ .collect(Collectors.toList());
+
+ if (uidTransitiveInEdges.size() < 2) {
+ return transitiveInEdgesInOrder;
+ }
+
+ int uidTransitiveInEdgeIndex = 0;
+ for (int transitiveInEdgeIndex = 0;
+ transitiveInEdgeIndex < transitiveInEdgesInOrder.size();
+ transitiveInEdgeIndex++) {
+ if
(hasUidBackedUpstream(transitiveInEdgesInOrder.get(transitiveInEdgeIndex))) {
+ transitiveInEdgesInOrder.set(
+ transitiveInEdgeIndex,
+ uidTransitiveInEdges.get(uidTransitiveInEdgeIndex++));
+ }
+ }
+ return transitiveInEdgesInOrder;
+ }
+
+ private boolean hasUidBackedUpstream(StreamEdge inEdge) {
+ return streamGraph
+ .getStreamNode(getStartNodeId(inEdge.getSourceId()))
+ .getTransformationUID()
+ != null;
+ }
+
+ private JobVertexID getStartNodeJobVertexId(
+ StreamEdge inEdge, JobVertexBuildContext jobVertexBuildContext) {
+ return new JobVertexID(
+ Preconditions.checkNotNull(
+
jobVertexBuildContext.getHash(getStartNodeId(inEdge.getSourceId()))));
+ }
Review Comment:
Would it make sense to deduplicate this logic across `AdaptiveGraphManager`
and `StreamingJobGraphGenerator`?
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]