使用Spark Streaming将Kafka Avro序列化数据存储为Parquet文件
Hey there! Let's put together this Kafka + Spark Streaming pipeline in Java, handling AVRO GenericRecord with Schema Registry and writing output to Parquet in append mode. I'll break this down into actionable steps with code examples you can adapt.
First, make sure your project includes all necessary dependencies. If you're using Maven, add these to your pom.xml (adjust versions to match your Spark/Kafka/Confluent setup):
<!-- Spark Streaming --> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming_2.12</artifactId> <version>3.3.0</version> <scope>provided</scope> </dependency> <!-- Kafka Integration for Spark Streaming --> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming-kafka-0-10_2.12</artifactId> <version>3.3.0</version> </dependency> <!-- Avro Support --> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-avro_2.12</artifactId> <version>3.3.0</version> </dependency> <!-- Confluent Schema Registry Client --> <dependency> <groupId>io.confluent</groupId> <artifactId>kafka-schema-registry-client</artifactId> <version>7.3.0</version> </dependency> <dependency> <groupId>io.confluent</groupId> <artifactId>kafka-avro-serializer</artifactId> <version>7.3.0</version> </dependency>
Next, initialize your Spark Streaming context and configure the Kafka consumer to deserialize AVRO GenericRecords using the Schema Registry. Here's how to set it up:
import org.apache.spark.SparkConf; import org.apache.spark.streaming.Duration; import org.apache.spark.streaming.api.java.JavaStreamingContext; import org.apache.spark.streaming.kafka010.ConsumerStrategies; import org.apache.spark.streaming.kafka010.KafkaUtils; import org.apache.spark.streaming.kafka010.LocationStrategies; import org.apache.avro.generic.GenericRecord; import org.apache.kafka.clients.consumer.ConsumerRecord; import java.util.Collections; import java.util.HashMap; import java.util.Map; public class KafkaAvroToParquetStream { public static void main(String[] args) throws InterruptedException { // 1. Initialize Spark Configuration SparkConf conf = new SparkConf() .setAppName("KafkaAvroToParquet") .setMaster("local[*]"); // Remove this line in production clusters // 2. Set up Streaming Context with 5-second batch interval JavaStreamingContext streamingContext = new JavaStreamingContext(conf, new Duration(5000)); // Enable checkpointing for fault tolerance (critical for append mode reliability) streamingContext.checkpoint("/path/to/checkpoint/directory"); // Use HDFS path in production // 3. Configure Kafka Consumer Parameters Map<String, Object> kafkaParams = new HashMap<>(); kafkaParams.put("bootstrap.servers", "your-kafka-broker:9092"); kafkaParams.put("group.id", "spark-streaming-consumer-group"); kafkaParams.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); // Use KafkaAvroDeserializer to read GenericRecord values kafkaParams.put("value.deserializer", "io.confluent.kafka.serializers.KafkaAvroDeserializer"); kafkaParams.put("schema.registry.url", "http://your-schema-registry:8081"); // Critical: Disable specific AVRO reader to work with GenericRecord kafkaParams.put("specific.avro.reader", "false"); // 4. Define the target Kafka topic String targetTopic = "your-avro-kafka-topic"; // 5. Create Direct Stream to consume from Kafka JavaInputDStream<ConsumerRecord<String, GenericRecord>> kafkaStream = KafkaUtils.createDirectStream( streamingContext, LocationStrategies.PreferConsistent(), ConsumerStrategies.Subscribe(Collections.singletonList(targetTopic), kafkaParams) ); // Extract GenericRecord values from Kafka records JavaDStream<GenericRecord> avroStream = kafkaStream.map(ConsumerRecord::value);
To write to Parquet, converting the GenericRecord stream to a Spark DataFrame is the most straightforward approach. We'll use foreachRDD to process each batch and handle the write:
// 6. Process each RDD and write to Parquet in append mode avroStream.foreachRDD((rdd, batchTime) -> { if (!rdd.isEmpty()) { // Reuse or initialize SparkSession SparkSession spark = SparkSession.builder() .config(rdd.sparkContext().getConf()) .getOrCreate(); // Convert RDD<GenericRecord> to DataFrame (using AVRO byte conversion) Dataset<Row> df = spark.read().format("avro").load( rdd.map(record -> { ByteArrayOutputStream out = new ByteArrayOutputStream(); DatumWriter<GenericRecord> writer = new GenericDatumWriter<>(record.getSchema()); Encoder encoder = EncoderFactory.get().binaryEncoder(out, null); writer.write(record, encoder); encoder.flush(); return RowFactory.create(out.toByteArray()); }) ); // Optional: Manually map fields if you need custom schema control // Dataset<Row> df = spark.createDataFrame(rdd.map(record -> { // return RowFactory.create( // record.get("user_id"), // record.get("event_timestamp"), // record.get("event_type") // ); // }), defineCustomSchema()); // 7. Write DataFrame to Parquet in append mode df.write() .mode(SaveMode.Append) .parquet("/path/to/parquet/output"); // Can be HDFS path like hdfs://cluster/storage/parquet } }); // Start the streaming job and wait for termination streamingContext.start(); streamingContext.awaitTermination(); } // Optional helper: Define custom Spark schema for manual mapping private static StructType defineCustomSchema() { return new StructType() .add("user_id", DataTypes.StringType) .add("event_timestamp", DataTypes.TimestampType) .add("event_type", DataTypes.StringType); } }
- Checkpointing: Always use a persistent HDFS path for checkpointing to ensure fault tolerance. This prevents data reprocessing after failures.
- Schema Compatibility: Ensure your AVRO schema aligns with Spark's Parquet schema. The automatic byte conversion method avoids manual schema mismatches.
- Resource Allocation: Remove
setMaster("local[*]")when deploying to a cluster, and configure executor cores/memory based on your workload. - Offset Management: Spark's Direct Stream manages Kafka offsets internally, but you can customize offset storage if needed.
内容的提问来源于stack exchange,提问作者Madhusudhan

