rohityadav1993 commented on code in PR #19067:
URL: https://github.com/apache/pinot/pull/19067#discussion_r3680646992


##########
pinot-connectors/pinot-spark-3-connector/src/main/scala/org/apache/pinot/connector/spark/v3/datasource/PinotDataWriter.scala:
##########
@@ -203,59 +206,163 @@ class PinotDataWriter[InternalRow](
     logger.info("Pushed segment tar file {} to: {}", (segmentTarFile.getName, 
destPath))
   }
 
+  /** Converts a Spark [[catalyst.InternalRow]] to a Pinot [[GenericRow]].
+   *
+   *  Spark's DataSourceV2 write path applies an `UnsafeProjection` 
immediately before invoking
+   *  `DataWriter.write(...)`, so every row reaching this method is an 
`UnsafeRow` whose array
+   *  fields are `UnsafeArrayData`. `UnsafeArrayData.array()` throws 
`UnsupportedOperationException`,
+   *  so the array branch must use accessors that work on both 
`UnsafeArrayData` and
+   *  `GenericArrayData` -- per-element iteration over the typed `getXxx(i)` 
methods.
+   *
+   *  Multi-value array fields are emitted as `Object[]` (boxed primitives, 
`String`s, or
+   *  `byte[]`s) because Pinot's segment-generation pipeline expects this 
shape. Stats collectors
+   *  in `pinot-segment-local` cast every MV entry as `(Object[])`, and 
`GenericRow.copy()`
+   *  clones array values via `(Object[]) value` -- a primitive Spark array 
(`int[]`, `long[]`,
+   *  ...) returned from `ArrayData.toIntArray()` / `toLongArray()` would 
`ClassCastException` at
+   *  the first stats pass and never produce a segment.
+   *
+   *  Nullability is handled at the field level: a top-level `isNullAt` guard 
funnels every null
+   *  field through `putValue(name, null)` so the scalar `StringType` branch 
never NPEs on a null
+   *  `getUTF8String(...)` and primitive branches never silently substitute 
`0` / `false` for a
+   *  null field value. Element-level nulls within a multi-value array are 
rejected with
+   *  `IllegalArgumentException`: Pinot's MV stats collectors NPE on null 
elements, so failing
+   *  fast in the writer with a clear message is strictly better than 
producing a corrupt
+   *  segment or surfacing an opaque downstream crash.
+   */
   private def internalRowToGenericRow(record: catalyst.InternalRow): 
GenericRow = {
     val gr = new GenericRow()
 
-    writeSchema.fields.zipWithIndex foreach { case(field, idx) =>
-      field.dataType match {
-        case org.apache.spark.sql.types.StringType =>
-          gr.putValue(field.name, record.getString(idx))
-        case org.apache.spark.sql.types.IntegerType =>
-          gr.putValue(field.name, record.getInt(idx))
-        case org.apache.spark.sql.types.LongType =>
-          gr.putValue(field.name, record.getLong(idx))
-        case org.apache.spark.sql.types.FloatType =>
-          gr.putValue(field.name, record.getFloat(idx))
-        case org.apache.spark.sql.types.DoubleType =>
-          gr.putValue(field.name, record.getDouble(idx))
-        case org.apache.spark.sql.types.BooleanType =>
-          gr.putValue(field.name, record.getBoolean(idx))
-        case org.apache.spark.sql.types.ByteType =>
-          gr.putValue(field.name, record.getByte(idx))
-        case org.apache.spark.sql.types.BinaryType =>
-          gr.putValue(field.name, record.getBinary(idx))
-        case org.apache.spark.sql.types.ShortType =>
-          gr.putValue(field.name, record.getShort(idx))
-        case org.apache.spark.sql.types.ArrayType(elementType, _) =>
-          elementType match {
-            case org.apache.spark.sql.types.StringType =>
-              gr.putValue(field.name, 
record.getArray(idx).array.map(_.asInstanceOf[String]))
-            case org.apache.spark.sql.types.IntegerType =>
-              gr.putValue(field.name, 
record.getArray(idx).array.map(_.asInstanceOf[Int]))
-            case org.apache.spark.sql.types.LongType =>
-              gr.putValue(field.name, 
record.getArray(idx).array.map(_.asInstanceOf[Long]))
-            case org.apache.spark.sql.types.FloatType =>
-              gr.putValue(field.name, 
record.getArray(idx).array.map(_.asInstanceOf[Float]))
-            case org.apache.spark.sql.types.DoubleType =>
-              gr.putValue(field.name, 
record.getArray(idx).array.map(_.asInstanceOf[Double]))
-            case org.apache.spark.sql.types.BooleanType =>
-              gr.putValue(field.name, 
record.getArray(idx).array.map(_.asInstanceOf[Boolean]))
-            case org.apache.spark.sql.types.ByteType =>
-              gr.putValue(field.name, 
record.getArray(idx).array.map(_.asInstanceOf[Byte]))
-            case org.apache.spark.sql.types.BinaryType =>
-              gr.putValue(field.name, 
record.getArray(idx).array.map(_.asInstanceOf[Array[Byte]]))
-            case org.apache.spark.sql.types.ShortType =>
-              gr.putValue(field.name, 
record.getArray(idx).array.map(_.asInstanceOf[Short]))
-            case _ =>
-              throw new UnsupportedOperationException(s"Unsupported data type: 
Array[${elementType}]")
-          }
-        case _ =>
-          throw new UnsupportedOperationException("Unsupported data type: " + 
field.dataType)
+    writeSchema.fields.zipWithIndex foreach { case (field, idx) =>
+      if (record.isNullAt(idx)) {
+        gr.putValue(field.name, null)
+      } else {
+        field.dataType match {
+          case StringType =>
+            gr.putValue(field.name, record.getUTF8String(idx).toString)
+          case IntegerType =>
+            gr.putValue(field.name, record.getInt(idx))
+          case LongType =>
+            gr.putValue(field.name, record.getLong(idx))
+          case FloatType =>
+            gr.putValue(field.name, record.getFloat(idx))
+          case DoubleType =>
+            gr.putValue(field.name, record.getDouble(idx))
+          case BooleanType =>
+            gr.putValue(field.name, record.getBoolean(idx))
+          case ByteType =>
+            gr.putValue(field.name, record.getByte(idx))
+          case BinaryType =>
+            gr.putValue(field.name, record.getBinary(idx))
+          case ShortType =>
+            gr.putValue(field.name, record.getShort(idx))
+          case ArrayType(elementType, _) =>
+            gr.putValue(field.name, convertArrayData(record.getArray(idx), 
field.name, elementType))
+          case _ =>
+            throw new UnsupportedOperationException("Unsupported data type: " 
+ field.dataType)
+        }
       }
     }
     gr
   }
 
+  private def convertArrayData(arr: ArrayData, columnName: String, 
elementType: DataType): AnyRef = {
+    val n = arr.numElements()
+    elementType match {
+      case StringType =>

Review Comment:
   can we make this a util function and reuse?
   gen AI susggets this:
   ```
   private def convertArrayData(arr: ArrayData, columnName: String, 
elementType: DataType): AnyRef = {
     val n = arr.numElements()
   
     def build[T <: AnyRef](get: Int => T): Array[T] = {
       val out = new Array[AnyRef](n).asInstanceOf[Array[T]]
       var i = 0
       while (i < n) {
         requireNoNullElement(arr, columnName, i)
         out(i) = get(i)
         i += 1
       }
       out
     }
   
     elementType match {
       case StringType  => build(i => arr.getUTF8String(i).toString)
       case BinaryType  => build(i => arr.getBinary(i))
       case IntegerType => build(i => Integer.valueOf(arr.getInt(i)))
       case LongType    => build(i => java.lang.Long.valueOf(arr.getLong(i)))
       case FloatType   => build(i => java.lang.Float.valueOf(arr.getFloat(i)))
       case DoubleType  => build(i => 
java.lang.Double.valueOf(arr.getDouble(i)))
       case BooleanType => build(i => 
java.lang.Boolean.valueOf(arr.getBoolean(i)))
       case ByteType    => build(i => java.lang.Byte.valueOf(arr.getByte(i)))
       case ShortType   => build(i => java.lang.Short.valueOf(arr.getShort(i)))
       case _ =>
         throw new UnsupportedOperationException(s"Unsupported data type: 
Array[$elementType]")
     }
   }
   ```



-- 
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