Skip to content

feat: admit string maps in Spark-to-Comet conversion - #6036

Merged
sunchao merged 2 commits into
apache:mainfrom
sunchao:codex/upstream-string-map-conversion-20260919
Oct 1, 2026
Merged

sunchao merged 2 commits into
apache:mainfrom
sunchao:codex/upstream-string-map-conversion-20260919

Conversation

@sunchao

@sunchao sunchao commented Sep 19, 2026

Copy link
Copy Markdown
Member

Which issue does this PR close?

Follow-up to #5954, extending the same Spark-to-Comet conversion boundary to MAP<STRING,STRING>.

Rationale for this change

The conversion gate still rejects every map type. Consequently, an explicitly enabled RDD or Parquet input containing a string map cannot feed otherwise supported native filters and projections. For example, a row-backed input with attributes MAP<STRING,STRING> falls back even when RDDScan conversion is enabled.

The existing Arrow map writer already supports this representation. Admit binary string keys and values through the conversion gate, including nullable maps and values and placement inside supported structs. Existing source opt-ins remain in effect.

What changes are included in this PR?

  • Add the narrow MapType(StringType, StringType, _) admission case.
  • Document supported nullability and explicit RDDScan configuration.
  • Cover both reader paths, retained output ownership, input reuse and slices, encoding-failure cleanup, nested structs, native shuffle round trips, and disabled-conversion controls. Other collection types and non-binary collations remain rejected.

How are these changes tested?

Local validation used Linux x86_64, JDK 21, Spark 4.1.3 and Scala 2.13.17, with a freshly built native debug library.

  • Complete CometArrowStreamSuite plus the SparkToColumnar tests in CometExecSuite: 48 passed, none failed or skipped.
  • Clean reactor test compilation, Spotless, Scalastyle, scoped Markdown formatting, and git diff --check: passed.

The map integration regression checks Spark result/schema parity and the actual conversion, native filter/projection, and native shuffle operators for RDD and Parquet inputs. Source conversion disabled controls retain fallback. Spark's broader SQL suite is requested through run-spark-4.1-tests.

AI assistance: adapted and validated with OpenAI Codex.

@sunchao sunchao added the run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue label Sep 19, 2026
@github-actions github-actions Bot added enhancement New feature or request area:ffi Arrow FFI / JNI boundary labels Sep 19, 2026

@rich7420 rich7420 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@sunchao thanks for the patch, LGTM

@viirya viirya left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the follow-up to #5954 — the shape of the change is right, and the CometArrowStreamSuite additions are unusually thorough (nonzero child offsets, zeroing the source buffers before asserting, ownership after closing the output, and an allocator-leak check on the encoding-failure path). I verified the core mechanics and they hold up: the JVM-side Arrow schema from Utils.toArrowField matches what native serde.rs builds (entries/key/value, sorted=false, non-null key), and the collation exclusion genuinely works.

My main question is about the shape of the gate rather than this line being wrong.

The narrow case contradicts a sibling gate over the same machinery. CometLocalTableScanExec overrides isTypeSupported only for NullType/interval types and otherwise falls through to DataTypeSupport's recursive rule, so it already admits arbitrary maps — including MAP<STRING,INT>, MAP<INT,STRING> and MAP<STRING,ARRAY<STRING>>, all of which this PR explicitly asserts are unsupported. Both operators mix in CometNativeArrowSource and convert through RowArrowReader → ArrowWriter, and ArrowWriter.createFieldWriter recurses into ListVector/MapVector children for every writer it knows. So the conversion machinery is not what's limiting us here.

That leaves two readings, and I can't tell which is intended: either the restriction is purely incremental, in which case dropping the override and inheriting the recursive rule would fix the inconsistency and pre-empt a long tail of one-line PRs; or some of those types are actually broken on this path, in which case CometLocalTableScanExec has a latent bug and the reason belongs in a code comment. Could you say which?

Test coverage regressed against the array test it mirrors. The array test loops v1 ∈ {"", format}, so it reaches the BatchScanExec (DSv2) branch of shouldApplySparkToColumnar. The map test never sets USE_V1_SOURCE_LIST, so it only exercises V1 FileSourceScanExec and RDDScanExec. That branch is separate in CometExecRule, and the existing "SparkToColumnar over BatchScan" test has no map either. JSON coverage was also dropped without explanation.

No duplicate-key coverage. spark.sql.mapKeyDedupPolicy=LAST_WIN lets duplicate keys survive into execution, and Arrow Map doesn't enforce uniqueness. Since the test already exercises partitionValues['hour'], it'd be worth confirming native map_extract picks the same entry Spark does.

Two nearly identical suite tests. The new test is close to line-for-line with the array one (same conf block, same conversions.size == 1, same shuffle assertions, same disabled-control). Worth parameterizing — and if the gate goes recursive, they collapse into one loop naturally.

Smaller things are inline. Nothing here is a correctness objection to the added line itself; I'd just like the design question answered before this lands, since each additional narrow case makes the eventual cleanup larger.

name: String,
fallbackReasons: ListBuffer[String]): Boolean = dt match {
case ArrayType(StringType, _) => true
case MapType(StringType, StringType, _) => true

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This line depends on a Scala/Spark subtlety that nothing here records: in Spark 4.x StringType is a case object extending StringType(UTF8_BINARY_COLLATION_ID, NoConstraint), and its equals compares collationId and constraint. So this pattern compiles to an object-equality check and a collated StringType instance does not match — which is exactly what makes the UTF8_LCASE assertions below pass.

The hazard is that "tidying" this to MapType(_: StringType, _: StringType, _) looks equivalent, compiles fine, and silently drops the collation guard. Same applies to the ArrayType(StringType, _) line above, inherited from #5954.

Could you add something like:

// `StringType` is the UTF8_BINARY case object; a collated StringType instance
// is not equal to it, so collated maps/arrays fall through to the reject case
// below. Do NOT rewrite these as `_: StringType`.

Also worth noting: the case _: ArrayType | _: MapType => false catch-all on the next line is order-dependent, and every future supported type has to be inserted above it. A brief note would help.

CometConf.COMET_SHUFFLE_MODE.key -> "native",
CometConf.COMET_CONVERT_FROM_PARQUET_ENABLED.key -> "true",
CometConf.COMET_SPARK_TO_ARROW_ENABLED.key -> "true",
CometConf.COMET_SPARK_TO_ARROW_SUPPORTED_OPERATOR_LIST.key -> "RDDScan") {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This conf block never sets SQLConf.USE_V1_SOURCE_LIST, so the Parquet cases here only reach the V1 FileSourceScanExec branch of shouldApplySparkToColumnar. The array test at line 3772 loops v1 ∈ {"", format} and therefore also covers the BatchScanExec (DSv2 ParquetScan) branch, which is a separate case in CometExecRule.

So the DSv2 admission path is currently untested for maps — the existing "SparkToColumnar over BatchScan" test doesn't include a map column either. Adding v1 to the loop the same way the array test does would close this. Same question for JSON: the array test covers it and this one drops it, and the PR description doesn't say whether that was deliberate.

}
}

test("SparkToColumnar string maps cross RDD and Parquet native boundaries") {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This test is close to line-for-line with "SparkToColumnar string arrays cross JSON and Parquet native boundaries" at line 3772 — same conf block, same conversions.size == 1 / supportsColumnar assertions, same native-shuffle check, same limit(1) early-stop, same disabled-control block. The differences are the schema, the row fixtures, and the projected expression.

Worth folding into one test parameterized over the collection type. If the gate becomes recursive per my top-level comment, this falls out naturally as one more entry in a dataType loop rather than a third copy.

val schema = StructType(Seq(StructField("tags", mapType, nullable)))
val arrowSchema = Utils.toArrowSchema(schema, "UTC")
val values: Seq[Seq[(String, String)]] = Seq(
Seq("" -> "", "é" -> "東京", "a\u0000b" -> "duplicate", "b" -> "duplicate"),

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Minor, but the "duplicate" labels read as if this row covers duplicate map keys — the keys here (a\u0000b, b) are distinct and it's the values that repeat. The array test's "dup", "dup" at line 365 is genuinely about repeated elements, so the parallel naming is misleading for maps.

Duplicate keys are the case actually worth having: mapKeyDedupPolicy=LAST_WIN lets them reach execution, Arrow Map doesn't enforce uniqueness, and the CometExecSuite test already does partitionValues['hour'] lookups. Adding a row with a repeated key and asserting the lookup matches Spark would cover the one map-specific semantic that arrays don't have.

Comet accelerates Iceberg scans of Parquet files. See the [Iceberg Guide] for more information.

[Iceberg Guide]: iceberg.md
[iceberg guide]: iceberg.md

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Unrelated to this PR — this is prettier normalizing the reference-link label to lowercase. Harmless (CommonMark matches labels case-insensitively, so [Iceberg Guide] on line 32 still resolves), but it's drive-by churn in a docs section the PR doesn't otherwise touch. Either drop it or mention it in the description so reviewers don't have to work out whether the link broke.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Spark-to-Comet conversion rejected all maps, preventing opted-in sources containing string maps from feeding native operators.
  • Design approach: Admit MapType(StringType, StringType, _) through the existing conversion path.
  • Correctness / compatibility analysis: Checked Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Map schema layout, nullability, string copying and lookup semantics agree. The singleton StringType pattern excludes non-default collations. Spark’s LAST_WIN deduplicates during construction, and native lookup preserves Spark’s first-match behavior for stored duplicates.
  • Key design decisions: The narrow gate reuses existing readers and writers without adding another abstraction. Conversion retains its per-entry copying cost and existing source opt-ins. No reproducible performance regression was identified.
  • Implementation sketch: One production admission case, documentation updates, schema-gate tests, RDD/Parquet integration tests, and Arrow ownership, slicing and failure-cleanup tests.
  • Behavioral changes worth calling out: Nullable string maps and values, including maps inside supported structs, can now reach eligible native operators. Other collection shapes remain rejected.
  • Suggested improvements: No introduced P1/P2 issues found within this review. Existing discussion does not establish an unresolved P1/P2 blocker.

Reviewed the full four-file diff from 5fdc96199685061b7c67ad28651b4c0a3dcd6541 to 11bc927a51d389b3187616d80bd0b03106ec7676. Confirmed the PR is not a draft. Read existing reviews and all five review threads. Routed skills: review-comet-pr and audit-comet-expression for downstream map semantics.

Exact-head CI: 43 successful checks, 29 skipped, no failures or pending checks. Both required-check aggregates passed. CI tested a merge commit whose tree exactly matches the reviewed head. Logs confirm 954 execution tests passed, including the new map tests, and native map-lookup tests passed. All seven Spark 4.1 SQL shards also passed.

Validation limits: git diff --check passed. A focused local native test attempt reached its 45-second limit while compiling dependencies, before executing tests. JVM tests were not rerun locally. Non-default Spark runtime suites, macOS and Iceberg were not exercised by this CI run. Project source remains unchanged.

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'd like to get viirya's question about the shape of the gate answered before this goes in, so I tried to get some data on it. I removed the whole collection override in CometSparkToColumnarExec.isTypeSupported locally, so it inherits the recursive rule in DataTypeSupport, and ran 16 array and map shapes through it on Spark 4.1 with main merged in. That covered maps with int, date, decimal(20,2), array and struct keys or values, and arrays of int, decimal(38,10), timestamp, binary, double, boolean, arrays and structs. Each shape went through an RDD, a Parquet row read and a Parquet vectorized read with native filters, projections and element access above the conversion, plus a native shuffle from the RDD, and every query matched Spark. I also tried the #4789 shapes, where non-null map values and array elements feed map_entries, slice and array_insert, through the cache and RDD paths, and those came out right too. The base rule still rejects collated strings at every nesting level.

The _: ArrayType | _: MapType => false case goes back to #1741, which only carried over the old structs-only support set, so as far as I can tell the narrow case isn't guarding against a known conversion bug. Could we drop the override instead? That would also settle the StringType matching comment and the order-dependent catch-all that viirya raised, and the new CometExecSuite test could become a loop over types rather than a second copy of the array test. If you'd rather keep this PR to string maps, I'm happy to open an issue for the rest.

Once this is rebased, two places on main will contradict it. docs/source/user-guide/latest/in-memory-cache.md says CometSparkToColumnarExec declines ArrayType and MapType, and the comment above the struct columns in CometInMemoryCacheBenchmark says the same. I added both in #5543 after this branch was cut, and they were already wrong for ARRAY<STRING> by then, so that part is on me. Could you update them here while you rebase?

With main merged, the new tests and the existing SparkToColumnar tests pass locally on Spark 4.1 and on 3.5 with Scala 2.12, and removing the new MapType case fails both new CometExecSuite tests. The DSv2 BatchScanExec and JSON paths viirya mentioned also work for string maps when I run them by hand, so covering them only needs the extra loop dimensions in the test.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Spark-to-Comet conversion rejected maps, preventing enabled sources containing string maps from feeding eligible native operators.
  • Design approach: Admit MapType(StringType, StringType, _) through the existing conversion machinery.
  • Correctness / compatibility analysis: Compared relevant Spark sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Map layout, nullability, copying and lookup behavior agree. The StringType singleton pattern excludes non-default collations. Spark’s LAST_WIN deduplicates during construction, while stored duplicate keys retain first-match lookup semantics in both engines.
  • Key design decisions: The narrow admission rule preserves source configuration and existing collection restrictions. It adds no abstraction or copying algorithm. Conversion retains the existing per-entry copying cost.
  • Implementation sketch: One production admission case, documentation updates, schema-gate assertions, RDD/Parquet integration coverage, and reader ownership, slicing and failure-cleanup tests.
  • Behavioral changes worth calling out: Nullable string maps and values, including maps inside supported structs, can now cross the conversion boundary and native shuffle.
  • Suggested improvements: No introduced P1/P2 issues found within this review. Existing discussion does not establish an unresolved P1/P2 blocker.

Reviewed the full four-file diff from 5fdc96199685061b7c67ad28651b4c0a3dcd6541 to 11bc927a51d389b3187616d80bd0b03106ec7676. Confirmed the PR is not a draft and read all existing reviews and five review threads. Routed skills: review-comet-pr and the relevant downstream-semantics checks from audit-comet-expression.

Exact-head CI: 43 successful checks, 29 skipped, no failures or pending checks. Both required-check aggregates passed. Verified that the CI merge commit has exactly the reviewed head’s tree. Logs confirm 954 execution tests passed, including the new map tests, and native map-lookup tests passed. All seven Spark 4.1 SQL shards passed.

Validation limits: git diff --check passed. JVM/native suites were not rebuilt or rerun locally; runtime validation relies on verified CI. Non-default Spark runtime suites, macOS and Iceberg were not exercised by these CI runs. Project source remains unchanged.

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM. Thanks @sunchao

@andygrove
andygrove enabled auto-merge September 30, 2026 12:37

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Approving, since nothing here is a correctness or performance problem. The new tests and the existing SparkToColumnar tests passed when I ran them with main merged, on 4.1 and on 3.5. Dropping the rest of the collection override and fixing the two stale passages from #5543 can come in a follow-up.

@andygrove
andygrove added this pull request to the merge queue Sep 30, 2026
@github-merge-queue
github-merge-queue Bot removed this pull request from the merge queue due to failed status checks Sep 30, 2026
@sunchao
sunchao added this pull request to the merge queue Oct 1, 2026
Merged via the queue into apache:main with commit 92b43cf Oct 1, 2026
49 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:ffi Arrow FFI / JNI boundary enhancement New feature or request run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants