loustler opened a new pull request, #11660:
URL: https://github.com/apache/seatunnel/pull/11660
`ParquetWriteStrategy#getOrCreateOutputStream` is called from `write()`,
i.e. once per row, and it constructs a `GenericData` and registers three
logical-type conversions on **every** call — but the object it builds is only
consumed on the branch that creates a new `ParquetWriter`. For every row that
writes into an already-open file it is built and immediately discarded.
I found this while profiling a production ingestion run, so there are
measurements attached. I want to be upfront about what they do and do not show:
**this is wasted CPU, wasted allocation and a contended JVM-global lock on the
write path, not a demonstrated end-to-end throughput win.** In my own run the
writers were parked 56.7% of the time waiting for rows, so the job was
source-bound and removing this cost would not have shortened it much. I've
tried to separate measured from inferred throughout, and I'd welcome correction
where I've misread something.
### The code
`seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/ParquetWriteStrategy.java`
```java
public void write(@NonNull SeaTunnelRow seaTunnelRow) { // once per
row
...
ParquetWriter<GenericRecord> writer = getOrCreateOutputStream(filePath);
```
```java
public ParquetWriter<GenericRecord> getOrCreateOutputStream(@NonNull String
filePath) {
...
ParquetWriter<GenericRecord> writer =
this.beingWrittenWriter.get(filePath);
GenericData dataModel = new GenericData();
// every row
dataModel.addLogicalTypeConversion(new Conversions.DecimalConversion());
// every row
dataModel.addLogicalTypeConversion(new
TimeConversions.DateConversion()); // every row
dataModel.addLogicalTypeConversion(new
TimeConversions.LocalTimestampMillisConversion());
if (writer == null) {
// only use site
...
AvroParquetWriter.<GenericRecord>builder(outputFile).withDataModel(dataModel)...
}
return writer;
// dataModel dropped
}
```
Those four lines have been unchanged since #2943 (Oct 2022) and are
identical at `2.3.11`, `2.3.13` and `dev`, so this is long-standing rather than
a regression.
### Why it is not just an allocation
The Avro that `connector-file-base` actually ships is **1.10.2**, via
`<parquet-avro.version>1.12.3</parquet-avro.version>`, relocated to
`org.apache.seatunnel.shade.connector.file.org.apache.avro`. Confirmed from the
built jar rather than from the dependency tree:
```
$ unzip -p connector-file-hadoop-2.3.13.jar
META-INF/maven/org.apache.avro/avro/pom.properties
version=1.10.2
$ unzip -l connector-file-hadoop-2.3.13.jar | grep generic/GenericData.class
org/apache/seatunnel/shade/connector/file/org/apache/avro/generic/GenericData.class
```
At that version `GenericData` has a field initialiser, which therefore runs
in every constructor
([`GenericData.java:182-183`](https://github.com/apache/avro/blob/release-1.10.2/lang/java/avro/src/main/java/org/apache/avro/generic/GenericData.java#L182-L183)):
```java
public static final String FAST_READER_PROP = "org.apache.avro.fastread";
private boolean fastReaderEnabled =
"true".equalsIgnoreCase(System.getProperty(FAST_READER_PROP));
```
On **JDK 8** — which is what the project's own image builds on
(`seatunnel-dist/src/main/docker/Dockerfile`: `FROM
seatunnelhub/openjdk:8u342`) and what my clusters run — `System.getProperty`
resolves to a `synchronized` method on the JVM-global system-properties object:
```
$ javap -c -cp $JAVA8/jre/lib/rt.jar java.util.Properties
public java.lang.String getProperty(java.lang.String);
2: invokespecial // Method
java/util/Hashtable.get:(Ljava/lang/Object;)Ljava/lang/Object;
$ javap -cp $JAVA8/jre/lib/rt.jar java.util.Hashtable | grep
'get(java.lang.Object)'
public synchronized V get(java.lang.Object);
```
So every row taken by every Parquet writer thread acquires the same JVM-wide
monitor. JDK 9 replaced `Properties`' backing store with a `ConcurrentHashMap`,
so the lock disappears on newer JDKs while the call cost remains — verified on
JDK 21: `getProperty` compiles to `ConcurrentHashMap.get`.
`addLogicalTypeConversion` then does, per call, a `HashMap.put` plus an
`IdentityHashMap` lookup and — on a freshly built model, which is every time
here — a `new LinkedHashMap<>()` per converted type
([`GenericData.java:126-135`](https://github.com/apache/avro/blob/release-1.10.2/lang/java/avro/src/main/java/org/apache/avro/generic/GenericData.java#L126-L135)).
### Evidence
**1. Production profile (measured; sampling, not instrumentation)**
SeaTunnel 2.3.13 Zeta, 20 workers, a 49-job DAG landing 30.45 billion rows
to HDFS Parquet+Snappy. JVM thread dumps via `GET /seatunnel/thread-dump` every
~6 s across all 20 workers, filtered to sink-writer threads, bucketed by
`frame0 | first non-JDK frame`. **96,078 writer samples, peak 198 concurrent
writer threads.**
| share | frame0 \| first non-JDK frame |
|---|---|
| 56.66% | `sun.misc.Unsafe.park` \| `MultiTableWriterRunnable.run` (idle,
queue empty) |
| **10.71%** | **`java.util.Hashtable.get` \|
`avro.generic.GenericData.<init>`** |
| 6.81% | `sun.misc.Unsafe.unpark` \| `MultiTableWriterRunnable.run` |
| 3.69% | `UTF_8$Encoder.encodeBufferLoop` \|
`parquet.io.api.Binary…encodeUTF8` |
| 3.53% | `ParquetWriteStrategy.write` |
| 1.94% | `SnappyNative.rawCompress` |
| 0.85% | `avro.generic.GenericData.<init>` |
| 0.59% | `HashMap.putVal` \| `GenericData.addLogicalTypeConversion` |
| 0.51% | `GenericData.addLogicalTypeConversion` |
All `GenericData`-attributed frames together: **12.99% of sink-writer
samples (12,479 / 96,078)**. Thread states over the same samples: RUNNABLE
46.8%, TIMED_WAITING 50.2%, BLOCKED 2.2%.
Reading this honestly: most `Hashtable.get` samples were RUNNABLE, i.e. the
threads mostly acquired the monitor uncontended. The dominant cost here is the
**work**, and lock contention is a smaller secondary effect at this concurrency
(~10 writer threads per JVM).
**2. Microbenchmark against the shaded Avro (measured)**
Those four lines only, compiled and run against the shipped
`connector-file-hadoop-2.3.13.jar` so the class under test is the relocated
one. Time-based rounds, 3 warmup + median of 5 × 1000 ms, per-thread blackhole
field so the allocation escapes as it does in the real method. 12-core machine,
`-Xms1g -Xmx1g`.
| allocation, 1 thread | JDK 8 | JDK 21 |
|---|---|---|
| the four lines as they are today | **1376 B/op** | 1456 B/op |
| `new GenericData()` only | 608 B/op | 664 B/op |
| `System.getProperty` only | 0 B/op | 0 B/op |
| a pre-built model handed back | **0 B/op** | **0 B/op** |
At my run's 30.45e9 rows that is ≈ **42 TB** of garbage produced and
immediately discarded.
| throughput, JDK 8 | threads | ops/s | ns/op | scaling |
|---|---|---|---|---|
| current | 1 | 8,712,249 | 114.8 | 1.00× |
| current | 8 | 7,104,190 | 1126.1 | 0.82× |
| current | 32 | 6,789,794 | 4713.0 | **0.78×** |
| `getProperty` only | 32 | 52,391,791 | 610.8 | **0.30×** |
| pre-built model | 32 | 2,213,908,086 | 14.5 | **5.46×** |
The current code's *aggregate* throughput does not improve at all as threads
are added — it caps at roughly 7M constructions/s per JVM regardless of core
count, and per-op latency grows linearly with thread count. On JDK 21 the same
code scales 4.82× at 32 threads, and `getProperty` alone 7.38×.
**3. The same cost on the real write path (measured)**
To check that this survives on the actual `write()` path — and in particular
that HotSpot is not eliminating the allocation — I drove
`ParquetWriteStrategy#write` directly: 200,000 pre-built rows per thread over a
9-column schema, timing the write loop only (`finishAndCloseFile()` closes one
writer per output file and dominates wall time; it is identical either way and
deliberately excluded). JDK 8, `-Xms3g -Xmx3g`, 21 rounds after a warmup round,
local filesystem. Patched and unpatched were alternated A/B/A/B across separate
JVMs so machine drift cannot favour one side.
Allocation per row is deterministic — min and max across 7 rounds land
within 1%:
| | unpatched | patched | delta |
|---|---|---|---|
| one output file (no rotation) | 2726.7 B/row | 1350.7 B/row | **−1376.0
B/row (−50.5%)** |
| `batch_size=2000` | 4817.5 B/row | 3474.2 B/row | −1343.3 B/row (−27.9%) |
The −1376.0 B/row matches the microbenchmark's 1376 B/op to the byte, which
settles the question: the object is **not** scalar-replaced on the real path,
and **half of everything the Parquet write path allocates is this
immediately-discarded object graph.**
Throughput, 8 writer threads, one output file, three independent A/B pairs —
the ranges do not overlap on either metric:
| | unpatched (3 runs) | patched (3 runs) |
|---|---|---|
| rows/s | 1.59 / 1.71 / 1.83 M | **2.15 / 2.30 / 2.31 M** |
| CPU ns/row | 2044 / 2103 / 2105 | **1489 / 1522 / 1621** |
**+17%** worst-patched vs best-unpatched, **+34%** median to median, **−21%
to −28%** CPU per row.
**Where the gain does not show up, stated plainly.** Only that one
configuration is resolvable.
- At **1 thread** the expected effect is ~115 ns against a ~2500 ns/row
baseline (≈5%), and this harness's JVM-to-JVM spread is larger than that; the
ranges overlap, so I am not quoting a number. One patched run even measured 7%
*slower*, which is not physically possible for strictly less work — that is the
noise floor, visible.
- With **frequent file rotation** (`batch_size=2000`) the wall-clock gain
disappears into noise (+0.7% at 1 thread, −7.3% at 8 threads). The reason is a
different per-file cost: each writer creation calls
`getConfiguration(hadoopConf)` → `hadoopConf.toConfiguration()` → `new
Configuration()`, re-parsing `core-default.xml` per file. Measured, ~48 ms per
file, which makes the per-row delta ~0.9% of wall — below what this harness can
resolve. CPU per row (−6.4%) and allocation (−24.6%) still improve there.
CPU per row and allocation per row improved in **every** configuration
measured; wall-clock throughput is only cleanly separable in the
multi-threaded, large-file case.
Caveats on this hardware: the only JDK 8 available to me is x86_64 under
Rosetta on an arm64 machine, so absolute values are not representative and only
ratios should be read; and the sink writes to a local filesystem, not HDFS,
where per-file network I/O would dilute the relative gain further.
**4. Reproducing the production frame, and isolating the lock (measured)**
The same sampling methodology against 10 threads executing only those four
lines reproduces the production bucket, and shows the full JDK path the
production dump elides:
```
94.66% java.util.Hashtable.get | avro.generic.GenericData.<init>
at java.util.Hashtable.get(Hashtable.java:363)
at java.util.Properties.getProperty(Properties.java:969)
at java.lang.System.getProperty(System.java:736)
at …avro.generic.GenericData.<init>(GenericData.java:183)
at …avro.generic.GenericData.<init>(GenericData.java:99)
```
| identical code and thread count | JDK 8 | JDK 21 |
|---|---|---|
| BLOCKED samples | **14,938 / 35,809 (41.7%)** | **4 / 45,840 (0.009%)** |
That is the lock, directly: same bytecode, same thread count, contention
only where `Properties` is a `synchronized Hashtable`.
### The change
```diff
ParquetWriter<GenericRecord> writer =
this.beingWrittenWriter.get(filePath);
- GenericData dataModel = new GenericData();
- dataModel.addLogicalTypeConversion(new
Conversions.DecimalConversion());
- dataModel.addLogicalTypeConversion(new
TimeConversions.DateConversion());
- dataModel.addLogicalTypeConversion(new
TimeConversions.LocalTimestampMillisConversion());
if (writer == null) {
+ GenericData dataModel = new GenericData();
+ dataModel.addLogicalTypeConversion(new
Conversions.DecimalConversion());
+ dataModel.addLogicalTypeConversion(new
TimeConversions.DateConversion());
+ dataModel.addLogicalTypeConversion(
+ new TimeConversions.LocalTimestampMillisConversion());
Path path = new Path(filePath);
```
The four statements move into the only branch that uses them.
`AbstractWriteStrategy` rotates to a new file every `batch_size` rows
(`currentBatchSize >= batchSize → newFilePart()`), so this turns a per-row cost
into a per-`batch_size`-rows cost — with `batch_size` in the low thousands, a
~2000× reduction in constructions.
**Why not hoist it to a field or a static**
- A field would buy nothing measurable. What is left after this change is
one construction per output file — ~115 ns against the ~48 ms that creating a
writer already costs (dominated by the `new Configuration()` noted above), i.e.
0.0002% of it.
- A `static` shared instance would be worse than it looks: `GenericData`
keeps a `Collections.synchronizedMap(new WeakHashMap<>())` default-value cache,
so sharing one instance across all writer threads would swap one JVM-global
monitor for another rather than removing one.
- Keeping it inside the branch is a strictly local change with no new shared
state and no lifetime questions.
An instance field would also be *safe*, for the record — one
`ParquetWriteStrategy` is only ever touched by one writer thread
(`MultiTableSinkWriter` gives each queue index its own writer set and its own
`MultiTableWriterRunnable`, and the commit paths `synchronized` on that
runnable), and the class is already not thread-safe (`beingWrittenWriter` is a
plain `LinkedHashMap`). It just isn't worth the extra surface. Happy to switch
if maintainers prefer it.
### Behaviour
The moved object is only ever read, never mutated, by the writer:
`AvroWriteSupport` in parquet-avro 1.12.3 touches the model only via
`getConversionByClass`, `getField` and `resolveUnion`, and
`GenericData#getConversionByClass` / `getConversionFor` at 1.10.2 are plain
`Map.get`. Nothing observable depends on when the object is created.
Verified:
- Full build and all Parquet unit tests on **JDK 8** (`1.8.0_502`):
`ParquetReadStrategyTest` (16), `ParquetWriteStrategyTest` (2),
`ParquetWriteStrategyEvolutionTest` (2), `ParquetTypeCoercionTest` (3) — **23
tests, 0 failures**.
- End-to-end write of 30,000 rows × 9 columns (`STRING, INT, BIGINT,
DECIMAL(20,4), DATE, TIMESTAMP, BOOLEAN, DOUBLE, BYTES`) at `batch_size=2000` →
15 Snappy Parquet files, read back and dumped as parquet schema + every decoded
record with its Java type. The dumps for patched and unpatched are
**identical** (SHA-256 `c9349aaf…`), and show the conversions in effect
(`c_decimal=0.0001(BigDecimal)`, `c_date=2020-01-02(LocalDate)`).
One note on method, since I tried the more obvious check first: comparing
the **bytes** of the output files does not work, and not because of this
change. parquet-mr serialises each column chunk's `Set<Encoding>` in `HashSet`
order, and `Enum.hashCode()` is the identity hash, so footer byte order varies
**between JVM runs of identical code** — two unpatched runs differ in 18 bytes
(2 per column × 9 columns), all `0x0A`↔`0x06` swaps in the footer. Hence the
semantic comparison above.
### Checked and deliberately not changed
- **`ParquetReadStrategy.java`** has the same four-line block, but it sits
in `read(FileSourceSplit, Collector)` with the record loop nested inside it, so
it is already once-per-split. No change needed.
- **Avro ≥ 1.11.3 would make this worse.** `GenericData` gained a
`loadConversions()` that runs `ServiceLoader.load(Conversion.class,
classLoader)` in the constructor; 1.10.2 and 1.11.1 have no `ServiceLoader`
reference at all (checked with `javap` on all four jars). So if `parquet-avro`
is ever bumped past that — `connector-iceberg` already pins Avro 1.11.3 — this
call site would pick up a classpath service lookup per row on top of everything
above. Fixing it now removes that future footgun. This is the one place where
I'd originally guessed wrong myself: I assumed the `Hashtable` was on the
`ServiceLoader` path before reading the source, and it is not.
- `write()` also calls `fieldName.toLowerCase()` per field per row (1.77% of
my writer samples). It looks precomputable, but it is a separate change and
I've left it out to keep this reviewable.
- The per-file `new Configuration()` re-parse found while measuring (~48 ms
per output file) is a much larger cost than this one, but it is unrelated and
belongs in its own issue.
--
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]