diff --git a/runners/flink/src/main/java/org/apache/beam/runners/flink/execution/JobResourceEnvironmentSession.java b/runners/flink/src/main/java/org/apache/beam/runners/flink/execution/JobResourceEnvironmentSession.java index 14065ba5dd4e3..538fa8eae1462 100644 --- a/runners/flink/src/main/java/org/apache/beam/runners/flink/execution/JobResourceEnvironmentSession.java +++ b/runners/flink/src/main/java/org/apache/beam/runners/flink/execution/JobResourceEnvironmentSession.java @@ -107,6 +107,8 @@ public void close() throws Exception { // TODO: eventually use this for reference counting open sessions. this.isClosed = true; client.close(); - dataServer.close(); + // TODO: Close data server when we solve bundle completion/channel close races between runners + // and SDK harnesses. + // dataServer.close(); } }