Empresas
Empregos
  • Sobre nós
  • Soluções
    • Publicação de vagas
      Publique sua vaga e receba candidatos qualificados em 48h.
    • Avaliações de candidatos
      Mais de 500 testes técnicos e psicológicos, mais anti-fraude.
    • Headhunting
      Busca executiva personalizada do início ao fim.
    • Folha de Pagamento + EOR
      Dispersão de folha e EOR em mais de 15 países da LATAM.
  • Preços
  • Empregos

0

355
Visualizações
Conjunto vacío después de collectAsList, aunque no esté vacío dentro del operador de transformación

Estoy tratando de averiguar si puedo trabajar con Kotlin y Spark, y usar las clases de datos del primero en lugar de las clases de casos de Scala.

Tengo la siguiente clase de datos:

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

Y la parte relevante de la rutina principal se ve así:

 val transactionEncoder = Encoders.bean(Transaction::class.java) val transactions = inputDataset .groupByKey(KeyExtractor(), KeyExtractor.getKeyEncoder()) .mapGroups(TransactionCreator(), transactionEncoder) .collectAsList() transactions.forEach { println("collected Transaction=$it") }

Con TransactionCreator definido como:

 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") } } }

Sin embargo, creo que me estoy encontrando con algún tipo de problema de serialización, porque el conjunto termina vacío después de la recolección. Veo el siguiente resultado:

 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=[])

Probé un KryoRegistrator personalizado para ver si era un problema con el HashSet de Kotlin:

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

Pero no parece ayudar. ¿Alguna otra idea?

Código completo aquí .

over 4 years ago · Santiago Trujillo
1 Respostas
Responde à pergunta

0

Parece ser un problema de serialización. La documentación de los estados Encoders.bean (Spark v2.4.0):

tipos de colección: solo matriz y java.util.List actualmente, el soporte de mapas está en progreso

Portar la clase de datos Transaction a Java y cambiar items a java.util.List parece ayudar.

over 4 years ago · Santiago Trujillo Relatório
Responde à pergunta
Encontrar trabalhos remotos

Descubra a nova forma de encontrar um emprego!

melhores empregos
Principais categorias de trabalho
Empresas
Postar vaga Preços Comercial
Jurídico
Termos e Condições Política de privacidade
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Recomende algumas ofertas para mim
Preciso de ajuda