andygrove opened a new issue, #5488: URL: https://github.com/apache/datafusion-comet/issues/5488
### Describe the bug `Utils.getFieldVector` accepts a fixed list of Arrow vector types that excludes `LargeVarCharVector` and `LargeVarBinaryVector`, throwing `Unsupported Arrow Vector for serialize` for either. Comet otherwise supports those representations deliberately: - `CometPlainVector` tracks whether the variable-width offsets are 64-bit, with the comment "PyArrow UDFs can hand back large_string / large_binary columns". - `Utils.toArrowType` maps `ArrowType.LargeUtf8` to `StringType`, again noting PyArrow UDF output. - `ArrowWriters` has `LargeStringWriter` and `LargeBinaryWriter`. So a batch carrying a large-offset column can be produced and read, but not serialized. On `main` the reachable consumers are the two `Utils.serializeBatches` call sites: - `org.apache.spark.sql.comet.operators.getByteArrayRdd` (`operators.scala`), used to collect a native plan's output as bytes. - `CometBroadcastExchangeExec.getByteArrayRdd`. The producer is `CometMapInBatchExec`, whose runner preserves the Arrow vectors a Python worker returns. A `mapInArrow` or PyArrow UDF returning `pa.large_string()` or `pa.large_binary()` whose result is then broadcast or collected should therefore fail at serialization. ### Steps to reproduce Not run end to end. A `mapInArrow` returning `pa.large_string()`, with the result on the build side of a broadcast join (or collected), is the shape expected to hit it. The rejection itself is trivially reproducible at the `Utils` level: build a `CometPlainVector` over a `LargeVarCharVector` holding two short strings, put it in a `ColumnarBatch`, and call `Utils.serializeBatches`. An ordinary `VarCharVector` succeeds on the same path. ### Expected behavior A batch Comet can produce and read should be serializable, or should be normalized to a representation that is, rather than failing at the serialization boundary. ### Additional context Raised during review of #5051 by @sunchao (https://github.com/apache/datafusion-comet/pull/5051#discussion_r3867601991), where it surfaced through the cache serializer. Filed separately because it is not specific to that PR and outlives it: `getFieldVector` and both `serializeBatches` consumers are on `main` today, while the cache serializer is not. What #5051 does and does not address, so this is not mistaken for fixed. It adds `Utils.isSupportedFieldVector`, a non-throwing predicate for the same vector list, and makes `Utils.isArrowBacked` answer false for a `CometVector` wrapping something `getFieldVector` rejects, so the cache write path converts such a batch instead of failing. That helper is a useful building block here, but it only reroutes the cache path: `getByteArrayRdd` and the broadcast path have no conversion fallback and still throw. If #5051 does not merge, both the helper and the predicate go with it. Two directions worth weighing: 1. Accept the large variants in `getFieldVector`. They are `FieldVector`s, so the question is whether everything downstream handles 64-bit offsets: Arrow IPC serialization, the decode path, and the native import in `ScanExec`. 2. Normalize 64-bit offsets to 32-bit before serialization, failing only when a column genuinely exceeds the 32-bit range. The first is cheaper if the downstream path is already clean; the second is safer and matches what the cache path now does by converting. -- 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]
