Empresas
Empleos
  • Sobre nosotros
  • Soluciones
    • Publicación de vacantes
      Publica tu vacante y recibe candidatos calificados en 48h.
    • Evaluación de candidatos
      500+ pruebas técnicas y psicológicas, más anti-fraude.
    • Headhunting
      Búsqueda ejecutiva a la medida de principio a fin.
    • Nómina + EOR
      Dispersión de nómina y EOR en más de 15 países de LATAM.
  • Precios
  • Empleos

0

357
Vistas
Empty set after collectAsList, even though it is not empty inside the transformation operator

I am trying to figure out if I can work with Kotlin and Spark, and use the former's data classes instead of Scala's case classes.

I have the following data class:

data class Transaction(var context: String = "", var epoch: Long = -1L, var items: HashSet<String> = HashSet()) :
    Serializable {
    companion object {
        @JvmStatic
        private val serialVersionUID = 1L
    }
}

And the relevant part of the main routine looks like this:

val transactionEncoder = Encoders.bean(Transaction::class.java)
val transactions = inputDataset
    .groupByKey(KeyExtractor(), KeyExtractor.getKeyEncoder())
    .mapGroups(TransactionCreator(), transactionEncoder)
    .collectAsList()

transactions.forEach { println("collected Transaction=$it") }

With TransactionCreator defined as:

class TransactionCreator : MapGroupsFunction<Tuple2<String, Timestamp>, Row, Transaction> {
    companion object {
        @JvmStatic
        private val serialVersionUID = 1L
    }

    override fun call(key: Tuple2<String, Timestamp>, values: MutableIterator<Row>): Transaction {
        val seq = generateSequence { if (values.hasNext()) values.next().getString(2) else null }
        val items = seq.toCollection(HashSet())
        return Transaction(key._1, key._2.time, items).also { println("inside call Transaction=$it") }
    }
}

However, I think I'm running into some sort of serialization problem, because the set ends up empty after collection. I see the following output:

inside call Transaction=Transaction(context=context1, epoch=1000, items=[c])
inside call Transaction=Transaction(context=context1, epoch=0, items=[a, b])
collected Transaction=Transaction(context=context1, epoch=0, items=[])
collected Transaction=Transaction(context=context1, epoch=1000, items=[])

I've tried a custom KryoRegistrator to see if it was a problem with Kotlin's HashSet:

class MyRegistrator : KryoRegistrator {
    override fun registerClasses(kryo: Kryo) {
        kryo.register(HashSet::class.java, JavaSerializer()) // kotlin's HashSet
    }
}

But it doesn't seem to help. Any other ideas?

Full code here.

over 4 years ago · Santiago Trujillo
1 Respuestas
Responde la pregunta

0

It does seem to be a serialization issue. The documentation of Encoders.bean states (Spark v2.4.0):

collection types: only array and java.util.List currently, map support is in progress

Porting the Transaction data class to Java and changing items to a java.util.List seems to help.

over 4 years ago · Santiago Trujillo Denunciar
Responde la pregunta
Encuentra empleos remotos

¡Descubre la nueva forma de encontrar empleo!

Top de empleos
Top categorías de empleo
Empresas
Publicar vacante Precios Comercial
Legal
Términos y condiciones Política de privacidad
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Recomiéndame algunas ofertas
Necesito ayuda