What is SerializeFromObject in spark execution plan?

Viewed 57

I read here that SerializeFromObject is performed when we go from DataSet to DataFrame. Is there official documentation that describes this in more detail?

Does this mean that going from DataSet to DataFrame is a relatively expensive operation, since it performs serialization?

1 Answers

Not sure about official documentation, but for Spark as open-source project, its source code would be the ultimate doc. And, according to SerializeFromObjectExec operator's source, it indeed does map an RDD of Java objects to an RDD of InternalRows using serializer (kryo-, java-, or other/custom).

/**
 * Takes the input object from child and turns in into unsafe row using the given serializer
 * expression.  The output of its child must be a single-field row containing the input object.
 */
case class SerializeFromObjectExec(
:
  override protected def doExecute(): RDD[InternalRow] = {
    child.execute().mapPartitionsWithIndexInternal { (index, iter) =>
      val projection = UnsafeProjection.create(serializer)
      projection.initialize(index)
      iter.map(projection)
    }
  }  
:

As we see here, the doExecute goes and applies UnsafeProjection (which "wraps" a serializer) on each partition to produce an RDD of InternalRows. In the query plan, we would also find a MapPartitions node added prior to SerializeFromObject.

Related