andygrove commented on code in PR #6412:
URL: https://github.com/apache/datafusion-comet/pull/6412#discussion_r4155982503
##########
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:
Added in 5c1e79cebe, `checkSizes` now runs over `withDictionaryCache` as
well. To make sure it catches the case you describe, I had the writer take a
dictionary column's size before decoding it. This test then fails with 2,063
bytes recorded for `s1` against 3,567 decoded, and nothing else in the suite
notices.
##########
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:
Done in 5c1e79cebe. The reasoning is now only at `statsRow`, citing
`ColumnStats.sizeInBytes` behind `DefaultCachedBatchSerializer`, and the
Scaladoc, `serialize`, the test helper and the two size tests point there. I
fixed the `encodeBatches` comment too. The same commit drops the suite's
`scala.collection.mutable` import, which was the only conflict between this PR
and #6421.
--
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]