ricky2129 commented on PR #12368:
URL: https://github.com/apache/seatunnel/pull/12368#issuecomment-5713701399
Thanks — Issue 1 is real and now fixed in `52aca0d31`.
I verified your finding against the pinned `parquet-avro:1.12.3` rather than
master, and it holds exactly as described. `ParquetReader.Builder#build()`
(1.12.3) does:
```java
if (path != null) {
FileSystem fs = path.getFileSystem(conf);
// allocation #1
FileStatus stat = fs.getFileStatus(path);
if (stat.isFile()) {
return new ParquetReader<>(
Collections.singletonList((InputFile)
HadoopInputFile.fromStatus(stat, conf)), // #2
...
```
and `HadoopInputFile.fromStatus` re-resolves via
`stat.getPath().getFileSystem(conf)` — so two uncached `FileSystem` instances
per file, neither closed, exactly the leak this PR set out to remove. It was
reachable deterministically: every file of any table whose column names aren't
Avro-safe goes through the `readWithAvro` -> `readWithNativeParquet` fallback.
My original audit grepped for `HadoopInputFile.fromPath` /
`HadoopOutputFile.fromPath` and therefore could not see this path, since the
resolution happens inside `build()`. Good catch.
**Fix approach.** 1.12.3 exposes no public `builder(ReadSupport,
InputFile)`, and `ParquetReader.Builder(InputFile)` is `protected`, so the
native path now uses a small private subclass that mirrors what
`AvroParquetReader.Builder` does internally:
```java
private static final class GroupParquetReaderBuilder extends
ParquetReader.Builder<Group> {
private GroupParquetReaderBuilder(InputFile file) { super(file); }
@Override protected ReadSupport<Group> getReadSupport() { return new
GroupReadSupport(); }
}
```
It also needs the explicit
`withConf(hadoopFileSystemProxy.getConfiguration())`, for the same reason you
identified on the Avro path — a non-`HadoopInputFile` falls into the `conf =
new Configuration()` branch — and it is applied before `withFileRange()` since
`withConf()` rebuilds the read options.
**Re-audited by mechanism rather than by factory name**, which is what let
the first one slip. Every site in `connector-file-base` that can resolve a
`FileSystem`:
| Site | State |
|---|---|
| `ParquetWriteStrategy` | fixed |
| `ParquetReadStrategy` Avro read (x2) | fixed |
| `ParquetReadStrategy` native/Group read | **fixed in `52aca0d31`** |
| `ParquetFileSplitStrategy:154` | fixed |
| `OrcReadStrategy` (x2) | fixed |
| `OrcWriteStrategy:160` | already correct on `dev` (`.fileSystem(...)`) |
| `ParquetFileSplitStrategy:146` (`proxy == null` branch) | not a leak —
builds a plain `new Configuration()` with the FS cache enabled, so Hadoop
returns a cached instance. Left unchanged deliberately. |
| `HadoopFileSystemProxy` `FileSystem.get` (x3) | the intended single
creation points |
**Issue 2 / test coverage.** Added
`testProxyOwnsASingleFileSystemSharedByManyFiles`, which asserts the proxy
hands back the same `FileSystem` across repeated lookups and across multiple
file writes/reads — a regression guard against any return to per-file
resolution. I did not add a `DefaultMetricsSystem` registration-count
assertion: that needs instrumentation of the static metrics system, and a naive
version on `file:///` would be misleading, since `HadoopConf`'s base schema is
`hdfs` and so the disable-cache key doesn't apply to the local FS in tests.
Happy to add it if you think it's worth the extra machinery.
`connector-file-base` on JDK 8: 40 tests, 0 failures. The only errors are 3
pre-existing `MockitoInitializationException` ("your JDK does not supply a
working agent attachment mechanism") in `ParquetFileSplitStrategyTest`, which
reproduce identically on unmodified `dev` in the same environment — unrelated
to this diff.
--
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]