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
AvroSparklerrecord wraps your actual data in a singleParsedDataKafkafield, which isn't necessary. We should model theParsedDataproperties 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
Metadatais essentially a key-value store—we don't need a nested record here; a simple Avromapwill work. - Unspecified
AnyReftype: Avro doesn't support arbitraryAnyRefvalues. We need to define explicit supported types for yourheadersmap.
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 uninitializedvar(which defaults tonull).outlinks: Avro doesn't have a nativeSettype, so we usearray<string>. We can convert between ScalaSetand Avroarrayduring serialization/deserialization.metadata: Maps directly to Tika'sMetadatausing an Avromap<string>. This assumes your metadata uses single-value keys; if you have multi-value keys, switch tomap<array<string>>.headers: The union value type covers common scalar types that fitAnyRef, 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-avroplugin oravro-maven-plugin) to generate a Scala/Java class from the corrected schema. You can then write a simple converter between this generated class and yourParsedDataclass. - Multi-Value Metadata: If your Tika
Metadatahas keys with multiple values, update the schema'smetadatafield 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
相关产品推荐
相关产品推荐

