Skip to content

Commit a214755

Browse files
committed
Address
1 parent 4466010 commit a214755

File tree

1 file changed

+1
-5
lines changed

1 file changed

+1
-5
lines changed

sql/core/src/main/scala/org/apache/spark/sql/execution/streaming/ForeachSink.scala

Lines changed: 1 addition & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -68,20 +68,16 @@ class ForeachSink[T : Encoder](writer: ForeachWriter[T]) extends Sink with Seria
6868
}
6969
datasetWithIncrementalExecution.foreachPartition { iter =>
7070
if (writer.open(TaskContext.getPartitionId(), batchId)) {
71-
var isFailed = false
7271
try {
7372
while (iter.hasNext) {
7473
writer.process(iter.next())
7574
}
7675
} catch {
7776
case e: Throwable =>
78-
isFailed = true
7977
writer.close(e)
8078
throw e
8179
}
82-
if (!isFailed) {
83-
writer.close(null)
84-
}
80+
writer.close(null)
8581
} else {
8682
writer.close(null)
8783
}

0 commit comments

Comments
 (0)