You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Avro新手求助:Schema与自定义ParsedData类适配问题

Fixing Avro Schema & Class Mapping for Your ParsedData Class

Hey there! Let's work through getting your ParsedData class properly mapped to an Avro schema—you're close, but there are a few tweaks needed to align the two properly. Let's break this down step by step.

First: Issues with Your Current Schema

Your existing schema has a few mismatches with your ParsedData class:

  • Redundant nesting: The top-level AvroSparkler record wraps your actual data in a single ParsedDataKafka field, which isn't necessary. We should model the ParsedData properties directly as top-level fields.
  • Missing field: Your class includes extractedText, but it's not defined in the schema.
  • Overly complex metadata definition: Tika's Metadata is essentially a key-value store—we don't need a nested record here; a simple Avro map will work.
  • Unspecified AnyRef type: Avro doesn't support arbitrary AnyRef values. We need to define explicit supported types for your headers map.

Corrected Avro Schema

Here's a schema that directly maps to your ParsedData class, with sensible defaults and type safety:

{
  "type": "record",
  "name": "ParsedData",
  "namespace": "edu.usc.irds.sparkler.avro",
  "fields": [
    {
      "name": "extractedText",
      "type": ["null", "string"],
      "default": null
    },
    {
      "name": "outlinks",
      "type": {
        "type": "array",
        "items": "string"
      },
      "default": []
    },
    {
      "name": "metadata",
      "type": {
        "type": "map",
        "values": "string"
      },
      "default": {}
    },
    {
      "name": "headers",
      "type": {
        "type": "map",
        "values": ["null", "string", "int", "long", "boolean"]
      },
      "default": {}
    }
  ]
}

Key Explanations for the Schema:

  • extractedText: Uses a union type (["null", "string"]) with a null default, matching your Scala class's uninitialized var (which defaults to null).
  • outlinks: Avro doesn't have a native Set type, so we use array<string>. We can convert between Scala Set and Avro array during serialization/deserialization.
  • metadata: Maps directly to Tika's Metadata using an Avro map<string>. This assumes your metadata uses single-value keys; if you have multi-value keys, switch to map<array<string>>.
  • headers: The union value type covers common scalar types that fit AnyRef, ensuring Avro can serialize/deserialize these values safely.

Mapping Between ParsedData and Avro Records

To convert between your custom class and Avro's GenericRecord (or generated classes), here's a Scala implementation:

Serialize ParsedData to Avro GenericRecord

import org.apache.avro.Schema
import org.apache.avro.generic.GenericData
import org.apache.avro.generic.GenericRecord
import org.apache.tika.metadata.Metadata
import scala.collection.JavaConverters._

def parsedDataToAvro(data: ParsedData, schema: Schema): GenericRecord = {
  val record = new GenericData.Record(schema)
  
  // Map extractedText
  record.put("extractedText", data.extractedText)
  
  // Convert Set[String] to Array[String] for Avro
  record.put("outlinks", data.outlinks.toArray)
  
  // Convert Tika Metadata to a Map[String, String]
  val metadataMap = data.metadata.names().asScala
    .map(name => name -> data.metadata.get(name))
    .toMap.asJava
  record.put("metadata", metadataMap)
  
  // Map headers, converting AnyRef to Avro-supported types
  val headersMap = data.headers.mapValues {
    case s: String => s
    case i: Int => i
    case l: Long => l
    case b: Boolean => b
    case _ => null // Handle unsupported types as null (or throw an error if preferred)
  }.asJava
  record.put("headers", headersMap)
  
  record
}

Deserialize Avro GenericRecord to ParsedData

def avroToParsedData(record: GenericRecord): ParsedData = {
  val data = new ParsedData()
  
  // Extract extractedText
  data.extractedText = record.get("extractedText").asInstanceOf[String]
  
  // Convert Array[String] back to Set[String]
  val outlinksArray = record.get("outlinks").asInstanceOf[Array[String]]
  data.outlinks = outlinksArray.toSet
  
  // Reconstruct Tika Metadata from the Avro map
  val metadataMap = record.get("metadata").asInstanceOf[java.util.Map[String, String]]
  val metadata = new Metadata()
  metadataMap.forEach((key, value) => metadata.add(key, value))
  data.metadata = metadata
  
  // Extract headers
  val headersMap = record.get("headers").asInstanceOf[java.util.Map[String, AnyRef]]
  data.headers = headersMap.asScala.toMap
  
  data
}

Additional Tips

  • Code Generation: For type safety, use Avro's code generation tools (like the sbt-avro plugin or avro-maven-plugin) to generate a Scala/Java class from the corrected schema. You can then write a simple converter between this generated class and your ParsedData class.
  • Multi-Value Metadata: If your Tika Metadata has keys with multiple values, update the schema's metadata field to use "type": "map", "values": {"type": "array", "items": "string"} and adjust the serialization/deserialization code to handle arrays.
  • Validation: Always validate your schema against your class fields to ensure no mismatches (missing fields, incorrect types) slip through.

内容的提问来源于stack exchange,提问作者Pilou Pierre-Louis Chevallier

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.27 09:23:23