ricky2129 opened a new pull request, #12368:
URL: https://github.com/apache/seatunnel/pull/12368
### Purpose of this pull request
The file connector lets Parquet and ORC resolve their own Hadoop
`FileSystem` instead of
using the one it already holds, and every such instance leaks its metrics
registration.
`HadoopConf#toConfiguration` sets `fs.<scheme>.impl.disable.cache = true`
unconditionally, so
`Path#getFileSystem(conf)` can never return a cached instance. Three APIs
rely on exactly
that call:
- `HadoopOutputFile.fromPath(path, conf)` — Parquet writes
- `HadoopInputFile.fromPath(path, conf)` — Parquet reads and footer reads
- `OrcFile.createReader(path, options)` when `ReaderOptions#filesystem` is
not set — ORC reads
So each file allocates a new `FileSystem` that nobody closes. The
`FileSystem` is then
garbage collected, but its constructor had registered an
`S3AInstrumentation` (or the
equivalent for other schemes) into the **static** `DefaultMetricsSystem`,
and only
`FileSystem#close()` unregisters it. The metrics objects therefore live for
the lifetime of
the JVM.
Heap histogram from a long-running worker writing Parquet to S3 (uptime 65.7
days, roughly
10 tables committing per 5-minute checkpoint):
| class | instances |
|---|---|
| `org.apache.hadoop.fs.s3a.S3AInstrumentation` | 222,061 |
| `org.apache.hadoop.metrics2.impl.MetricsSourceAdapter` | 222,062 |
| `org.apache.hadoop.metrics2.lib.MutableQuantiles` | 444,032 |
| `ScheduledThreadPoolExecutor$ScheduledFutureTask` | 444,067 |
| **`org.apache.hadoop.fs.s3a.S3AFileSystem`** | **71** |
222,061 instrumentations against 71 live `S3AFileSystem` instances is the
leak: the
filesystems are collected, their registrations are not. Each
`MutableQuantiles` schedules a
periodic rollover task, so a single metrics thread ends up servicing ~444k
timers. Observed
over 50 days at flat input volume: worker CPU 18% → 100%, heap growing
monotonically while
`major.gc.count` stayed 0, because nothing is reclaimable. A worker on the
same image with no
file sink was flat across the same window.
`OrcWriteStrategy` was already correct — it passes
`.fileSystem(hadoopFileSystemProxy.getFileSystem())` to
`OrcFile.writerOptions`. This makes
the remaining paths consistent with it.
### Changes
- add `FileSystemOutputFile` / `FileSystemInputFile`: Parquet
`OutputFile`/`InputFile`
adapters bound to a supplied `FileSystem`. `HadoopOutputFile` and
`HadoopInputFile` both
have private constructors and cannot be built from an existing
`FileSystem`. Stream
creation mirrors them exactly — same buffer size, `getDefaultReplication`,
`max(getDefaultBlockSize, blockSizeHint)`, `makeQualified`, and the same
`{hdfs, webhdfs, viewfs}` set for `supportsBlockSize()` — so behaviour is
unchanged for
every scheme.
- use them in `ParquetWriteStrategy`, `ParquetReadStrategy` (2 call sites)
and
`ParquetFileSplitStrategy`.
- set `ReaderOptions#filesystem` at both `OrcFile.createReader` call sites in
`OrcReadStrategy`.
- `ParquetReader.Builder` lifts a `Configuration` automatically only from a
concrete
`HadoopInputFile`, so the Parquet read paths now pass it explicitly via
`withConf(...)` and
`HadoopReadOptions`. Without this, `READ_INT96_AS_FIXED` and
`ADD_LIST_ELEMENT_RECORDS` set
by `HadoopConf` are silently dropped, which surfaces as
`IllegalArgumentException: INT96 is deprecated` and `ArrayStoreException`
on affected files.
`HadoopFileSystemProxy#getConfiguration` was added for this. Note
`withConf()` rebuilds the
read options, so it is applied before `withFileRange()`.
- the `hadoopFileSystemProxy == null` branch of `ParquetFileSplitStrategy`
is left unchanged:
it builds a plain `new Configuration()` with the cache enabled, so it does
not leak.
- drop a redundant second `getConfiguration(hadoopConf)` call in the Parquet
write path, which
constructed a whole `Configuration` per file.
### Does this PR introduce _user-facing_ change?
No. The same files are written and read with the same options; only the
`FileSystem` instance
used differs.
### How was this patch tested?
New `FileSystemParquetFileTest` covers a write/read round trip through both
adapters,
`createOrOverwrite` replacing an existing file, and that
`supportsBlockSize()`,
`defaultBlockSize()` and `getPath()` agree with the underlying `FileSystem`.
Existing suites pass on JDK 8: `ParquetReadStrategyTest` (16),
`ParquetWriteStrategyTest` (3),
`ParquetWriteStrategyEvolutionTest` (2), `ParquetTypeCoercionTest` (3),
`OrcReadStrategyTest`
(4), `OrcWriteStrategyTest` (4).
--
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]