comphead opened a new pull request, #6355:
URL: https://github.com/apache/datafusion-comet/pull/6355
## Which issue does this PR close?
Closes #6007.
## Rationale for this change
With `spark.comet.parquet.write.enabled=true`, both native write serdes,
`CometWriteFiles` on Spark 4.0+ and `CometDataWritingCommand` on Spark 3.x,
decline every output path that is not `file:` or `hdfs:`. Native Parquet writes
therefore never run against S3, where most production writes go.
Admitting `s3a:` alone would not be safe, for three reasons found on `main`:
- `objectstore::s3::create_store` hardcodes `AccessMode::Read`, and the
process-wide `object_store` cache has no access mode in its key. A write would
reuse a scan's store and present whatever a `CometS3CredentialProvider` grants
for reads.
- The writer receives Hadoop's unescaped `Path.toString` but derives the
object key the way the scan does, by percent-decoding the URL path. A name
holding `%`, `?` or `#` would be written under a different key from the one
Spark commits, and the job would succeed without its data.
- S3A does work at `create()` that an `object_store` upload bypasses: the
magic committer's pending uploads, and the encryption, ACL, storage class,
content encoding and headers it applies to every object it creates.
## What changes are included in this PR?
Native:
- `parquet_writer.rs`: `ParquetWriter::ObjectStore` wraps an
`AsyncArrowWriter` over `object_store::buffered::BufWriter`. A file under 10
MiB goes up in one PUT, a larger one as a multipart upload that only becomes an
object when `close` completes it. A write that fails before `close` aborts the
upload. The `bytes_written` metric now reports the uploaded size. It read 0 for
anything but a local file. URLs that normalize to `s3` (`s3`, `s3a` and
configured aliases) route to the new writer, and other object store schemes
still fail with `Unsupported storage scheme`.
- `parquet_support.rs`: `prepare_object_store_with_configs` is split into
`resolve_object_store` (normalization, backend selection and the cache) and the
`RuntimeEnv` registration. The new `object_store_for_write` resolves under
`AccessMode::Write` and registers nothing. The cache key gains the access mode.
Scans resolve and register exactly as before.
- `objectstore/s3.rs`: `create_store` takes the access mode it forwards to
the `CometS3CredentialBridge`.
- `native/core/Cargo.toml` names the `parquet` crate's `async` and
`object_store` features, which DataFusion already enables.
JVM:
- `NativeWriteUtils.unsupportedDestination` is now the filesystem gate of
both serdes. Local and HDFS destinations are gated as before. An S3 destination
is admitted only when:
- its scheme is `s3a` or `s3`, and Hadoop serves it with
`org.apache.hadoop.fs.s3a.S3AFileSystem` itself. `s3://` on EMRFS, subclasses
of S3A and the `fs.comet.s3Compliant.schemes` aliases stay on Spark.
- `fs.comet.libhdfs.schemes` does not route it through libhdfs.
- neither the path nor `mapreduce.output.basename` contains `%`, `?`, `#`
or an ASCII control character. Spaces and non-ASCII names survive the native
URL round trip and stay native.
- it is not under a `__magic*` directory, and `fs.s3a.committer.name`
(global or per-bucket) is not `magic`.
- none of `fs.s3a.encryption.algorithm`,
`fs.s3a.server-side-encryption-algorithm`, `fs.s3a.acl.default`,
`fs.s3a.create.storage.class`, `fs.s3a.object.content.encoding` or
`fs.s3a.create.header.*` is set for the bucket.
- `checkNativeWriteDestination`, which runs on the path the commit protocol
hands each task on Spark 4.0+, repeats the path and magic-directory checks.
- Docs: `operators.md`, `s3-credential-providers.md` and
`s3-credential-provider-design.md`, the `operator.proto` comment that counted
the S3A magic committer among the committers that work, and the
`CometS3AccessMode.WRITE` Javadoc.
Left out on purpose, possible follow-ups:
- Translating S3A's SSE settings into `object_store`'s encryption options
instead of falling back.
- `gs://`, `abfss://` and the `fs.comet.s3Compliant.schemes` aliases.
- Aborting the upload when a task is killed rather than failed. The writer
is then dropped without an abort, and the incomplete upload is left for the
bucket's lifecycle rules.
- Accounting for the upload buffers (10 MiB parts, up to 8 in flight per
writer) in Comet's memory pool.
## How are these changes tested?
- Rust unit tests, against an in-memory store seeded into the `object_store`
cache so that no request leaves the process:
- a small file lands in one PUT at the key the scan reads, through a
directory with a space and a non-ASCII name, and `bytes_written` matches its
size,
- a row group of about 16 MiB is streamed as a multipart upload and
completed,
- a write that fails after its multipart upload started aborts it and
leaves no object,
- other object store schemes are rejected before any store is built,
- a writer and a scan of the same bucket and configuration get separate
stores.
- `CometParquetWriterSuite`: `S3 destinations are admitted only where the
native writer matches S3A` covers each rule above without an S3 endpoint, plus
the task-time check and its agreement with planning.
- `ParquetWriteToS3Suite` is new and runs against MinIO. It is a manual
suite like the other MinIO suites, so it is listed in `dev/ci/check-suites.py`.
It checks that native writes are committed where Spark reads them (and
`_SUCCESS` on 4.0+), that a file larger than the upload buffer is uploaded in
parts (from its ETag), `Append` and `Overwrite`, that a directory with a space
and a non-ASCII name stays native, that `%` and the magic committer fall back,
and that a per-bucket `CometS3CredentialProvider` is asked for `WRITE` only.
Verified locally on the default Spark 4.1 profile: `cargo clippy
--all-targets --workspace -- -D warnings`, `cargo fmt --check`, and `./mvnw
test-compile` with scalastyle. The new tests have not been run yet, and neither
have the other Spark profiles or scalafix. `ParquetWriteToS3Suite` uses
`CometS3TestBase` and its `minio/minio:latest` image, like the other MinIO
suites, and anonymous pulls of MinIO images have been reported failing
(apache/datafusion#25215).
--
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]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]