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]

Reply via email to