comphead commented on code in PR #6412:
URL: https://github.com/apache/datafusion-comet/pull/6412#discussion_r4146860368
##########
spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala:
##########
@@ -1998,32 +1997,74 @@ class CometInMemoryCacheSuite extends CometTestBase {
}
}
- test("Comet in-memory cache records per-column sizes in its statistics") {
- // SimpleMetricsCachedBatch reserves a fifth field per column for its
size. A column owns a
- // known run of buffers in the payload, so the real stored size is known
and must be reported
- // rather than left at zero. Run over the nested relation as well: a
nested column's size is
- // the sum of its whole subtree, so this is also where a size attributed
to the wrong column
- // surfaces.
+ test("Comet in-memory cache records per-column decoded sizes in its
statistics") {
+ // SimpleMetricsCachedBatch reserves a fifth field per column for its size
and sums those into
+ // the batch's sizeInBytes, which is what Spark's planner reads as the
size of a materialized
+ // cached relation. Spark's own formats record a column's decoded size
there, so each field is
+ // compared with its column decoded back out of the payload, not with what
the column occupies
+ // compressed. Run over the nested relation as well: a nested column's
size is the sum of its
+ // whole subtree, so this is also where a size attributed to the wrong
column surfaces.
def checkSizes(
relation: org.apache.spark.sql.execution.columnar.InMemoryRelation,
batches: Array[CachedBatch]): Unit = {
val cacheSchema = Utils.fromAttributes(relation.output)
- batches.foreach { batch =>
- val sizes = CometCachedBatchHelper.columnSizes(batch, cacheSchema)
- val stats = batch.asInstanceOf[SimpleMetricsCachedBatch].stats
- sizes.zipWithIndex.foreach { case (size, i) =>
- assert(
- stats.getLong(i * 5 + 4) == size,
- s"column ${relation.output(i).name} should report the stored size
of its own " +
- "buffers in the statistics row")
+ val allocator = CometArrowAllocator.newChildAllocator("decoded-sizes",
0, Long.MaxValue)
+ try {
+ batches.foreach { batch =>
+ val sizes = CometCachedBatchHelper.decodedColumnSizes(batch,
cacheSchema, allocator)
+ val stats = batch.asInstanceOf[SimpleMetricsCachedBatch].stats
+ sizes.zipWithIndex.foreach { case (size, i) =>
+ assert(
+ stats.getLong(i * 5 + 4) == size,
+ s"column ${relation.output(i).name} should report its decoded
size in the " +
+ "statistics row")
+ }
+ assert(batch.sizeInBytes == sizes.sum)
}
+ } finally {
+ allocator.close()
}
}
withProjectionCache(checkSizes _)
withNestedProjectionCache(checkSizes _)
Review Comment:
Would it make sense to run `checkSizes` over `withDictionaryCache` too? That
fixture's shuffle hands the writer dictionary-encoded `s1` and `s2`, and
`serialize` decodes them before the size is taken. It is the one path where
measuring at the wrong point would understate the size again, and the
`range`-based fixtures here never reach `decodeDictionaries`. Something like
`withDictionaryCache { r => checkSizes(r,
r.cacheBuilder.cachedColumnBuffers.collect()) }` might work, though I haven't
run it.
##########
spark/src/main/scala/org/apache/spark/sql/comet/execution/arrow/ArrowCachedBatchSerializer.scala:
##########
@@ -350,11 +354,13 @@ class ArrowCachedBatchSerializer extends
SimpleMetricsCachedBatchSerializer {
values(base + 1) = upper(c)
values(base + 2) = nulls(c)
values(base + 3) = numRows
- // The stored size of the column's own Arrow buffers, taken from the
message's buffer
- // layout, so it is exact rather than an estimate. Cache pruning uses
- // bounds/null-count/row-count rather than this field, but Spark
reserves it and reports it,
- // so record the real value. The per-batch message framing is not
attributed to any column,
- // so these sum to slightly less than sizeInBytes.
+ // The column's decoded size: the plain length of its own Arrow buffers
before compression,
Review Comment:
Nit: this rationale now appears in the `CometCachedBatch` Scaladoc, here, in
`CachedBatchIpc.serialize`, and in three tests. Would it be enough to keep it
in one place? That one could also cite something present on every supported
Spark, such as the uncompressed `ColumnStats.sizeInBytes` behind
`DefaultCachedBatchSerializer`, because Spark's Arrow cache format is not in
3.5 to 4.1. Separately, the comment in `encodeBatches` above
`gatherColumnStats` still says the sizes are what "the message reports", but
they are now measured on the plain batch before compression.
--
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]