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

使用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.

Step 1: Set Up Dependencies

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>
Step 2: Read AVRO Data from Kafka

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);
Step 3: Convert AVRO to DataFrame & Write to Parquet

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);
    }
}
Key Tips for Production
  • 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 10:05:11