Open AllenWee1106 opened 10 months ago
This issue has been automatically marked as stale because it has been open for 180 days with no activity. It will be closed in next 14 days if no further activity occurs. To permanently prevent this issue from being considered stale, add the label 'not-stale', but commenting on the issue is preferred when possible.
Query engine
flink 1.16.2 iceberg 1.3.1
Question
flink 1.16.2 iceberg 1.3.1
tabenv.executeSql("create catalog jdbc with " + "('type'='iceberg'," + "'catalog-impl'='org.apache.iceberg.jdbc.JdbcCatalog'," + "'io-impl'='org.apache.iceberg.aws.s3.S3FileIO'," + "'uri'='jdbc:postgresql://***:****/jdbc'," + "'jdbc.user'='******'," + "'jdbc.password'='*****'," + "'warehouse'='s3://warehouse/'," + "'clients'='5'," + "'s3.endpoint'='http://***:***'," + "'s3.access-key-id'='******'," + "'s3.secret-access-key'='****'," + "'client.region'='cn-east-1'," + "'s3.path-style-access'='true'" + ")");
Table tableResult = tabenv.sqlQuery( "select * from jdbc.datalake.ods_vtc_tracking " + "where datetime >= '2023-11-07' " + "and datetime <= '2023-11-08'");
When executing the above SQL query, the following exception:
2023-12-06 10:32:45 WARN S3InputStream:237 - Unclosed input stream created by: org.apache.iceberg.aws.s3.S3InputStream.(S3InputStream.java:74)
org.apache.iceberg.aws.s3.S3InputFile.newStream(S3InputFile.java:85)
org.apache.iceberg.orc.FileIOFSUtil$InputFileSystem.open(FileIOFSUtil.java:125)
org.apache.iceberg.shaded.org.apache.orc.impl.ReaderImpl.extractFileTail(ReaderImpl.java:779)
org.apache.iceberg.shaded.org.apache.orc.impl.ReaderImpl.(ReaderImpl.java:567)
org.apache.iceberg.shaded.org.apache.orc.OrcFile.createReader(OrcFile.java:385)
org.apache.iceberg.orc.ORC.newFileReader(ORC.java:780)
org.apache.iceberg.orc.OrcIterable.iterator(OrcIterable.java:83)
org.apache.iceberg.io.CloseableIterable$ConcatCloseableIterable$ConcatCloseableIterator.hasNext(CloseableIterable.java:257)
org.apache.iceberg.io.CloseableIterable$7$1.hasNext(CloseableIterable.java:197)
org.apache.iceberg.io.CloseableIterable$7$1.hasNext(CloseableIterable.java:197)
org.apache.iceberg.relocated.com.google.common.collect.Iterators.addAll(Iterators.java:366)
org.apache.iceberg.relocated.com.google.common.collect.Iterables.addAll(Iterables.java:337)
org.apache.iceberg.deletes.Deletes.toEqualitySet(Deletes.java:116)
org.apache.iceberg.data.DeleteFilter.applyEqDeletes(DeleteFilter.java:192)
org.apache.iceberg.data.DeleteFilter.applyEqDeletes(DeleteFilter.java:211)
org.apache.iceberg.data.DeleteFilter.filter(DeleteFilter.java:154)
org.apache.iceberg.flink.source.RowDataFileScanTaskReader.open(RowDataFileScanTaskReader.java:97)
org.apache.iceberg.flink.source.DataIterator.openTaskIterator(DataIterator.java:139)
org.apache.iceberg.flink.source.DataIterator.updateCurrentIterator(DataIterator.java:129)
org.apache.iceberg.flink.source.DataIterator.hasNext(DataIterator.java:109)
org.apache.iceberg.flink.source.FlinkInputFormat.reachedEnd(FlinkInputFormat.java:125)
org.apache.flink.streaming.api.functions.source.InputFormatSourceFunction.run(InputFormatSourceFunction.java:89)
org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:110)
org.apache.flink.streaming.api.operators.StreamSource.run(StreamSource.java:67)
org.apache.flink.streaming.runtime.tasks.SourceStreamTask$LegacySourceFunctionThread.run(SourceStreamTask.java:333)