{"record":{"id":"e6386194fd4dd444","repo":"apache/seatunnel","slug":"close-file-output-stream-failed","errorCode":null,"errorMessage":"Close file output stream {} failed","messagePattern":"Close file output stream (.+?) failed","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/DebeziumJsonWriteStrategy.java","lineNumber":115,"sourceCode":"        }\n    }\n\n    @Override\n    public void finishAndCloseFile() {\n        beingWrittenOutputStream.forEach(\n                (key, value) -> {\n                    try {\n                        value.flush();\n                    } catch (IOException e) {\n                        throw new FileConnectorException(\n                                CommonErrorCodeDeprecated.FLUSH_DATA_FAILED,\n                                String.format(\"Flush data to this file [%s] failed\", key),\n                                e);\n                    } finally {\n                        try {\n                            value.close();\n                        } catch (IOException e) {\n                            log.warn(\"Close file output stream {} failed\", key, e);\n                        }\n                    }\n                    needMoveFiles.put(key, getTargetLocation(key));\n                });\n        beingWrittenOutputStream.clear();\n        isFirstWrite.clear();\n    }\n\n    @Override\n    public FSDataOutputStream getOrCreateOutputStream(@NonNull String filePath) {\n        FSDataOutputStream fsDataOutputStream = beingWrittenOutputStream.get(filePath);\n        if (fsDataOutputStream == null) {\n            try {\n                switch (compressFormat) {\n                    case LZO:\n                        LzopCodec lzo = new LzopCodec();\n                        OutputStream out =\n                                lzo.createOutputStream(","sourceCodeStart":97,"sourceCodeEnd":133,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/DebeziumJsonWriteStrategy.java#L97-L133","documentation":"Logged in DebeziumJsonWriteStrategy.finishAndCloseFile when closing an output stream for a finished file throws an IOException. Data was already flushed (flush failure is reported separately); this warning indicates the stream close itself failed, e.g. because the underlying Hadoop FS connection is broken or the file was deleted/moved concurrently. The file is still registered in needMoveFiles for commit.","triggerScenarios":"During sink close/commit, finishAndCloseFile iterates beingWrittenOutputStream and calls value.close(); the FSDataOutputStream.close throws IOException due to a failed HDFS/DataNode connection, deleted staging file, or quota/permission issue surfacing at close time.","commonSituations":"HDFS DataNode unavailability or network partitions at job end; distributed FS (S3/OSS) transient 5xx on finalize; file removed by a cleanup job before close; Kerberos token expiry on long jobs.","solutions":["Inspect the attached stack trace for the underlying IOException cause and fix the storage/network issue","Retry the job; SeaTunnel checkpoint/restart will rewrite the affected files","Verify the staging directory is not concurrently cleaned and the user has write permissions","Check storage backend health (DataNodes, S3 endpoint) and increase timeouts/retry settings"],"exampleFix":null,"handlingStrategy":"retry","validationCode":"// check storage reachability and staging dir writability before the job\nFileSystem fs = FileSystem.get(conf);\nPath staging = new Path(stagingDir);\nif (!fs.exists(staging) || !fs.getFileStatus(staging).getPermission().getAction().contains(org.apache.hadoop.fs.permission.FsAction.WRITE)) { throw new IllegalStateException(\"staging dir not writable: \" + staging); }","typeGuard":null,"tryCatchPattern":"// rely on framework retry; if wrapping, log-and-continue with verification\ntry {\n    strategy.finishAndCloseFile(...);\n} catch (Exception e) {\n    log.warn(\"File commit issue, will retry job/verify files\", e);\n    // verify target files exist and sizes > 0 before declaring success\n}","preventionTips":["Monitor HDFS/DataNode or object-store health during job windows","Keep jobs shorter than token/session expiry or enable token renewal","Never run external cleanup against the sink staging directory","Verify final file sizes/counts after jobs that logged close failures"],"tags":["ioexception","hadoop","file-close","storage"],"backgroundTag":"file-write-failed","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}