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

Spark Streaming消费Kafka Avro数据写入Parquet与Cassandra的Java实现咨询

Java实现:从Kafka Avro消息生成DataFrame并写入Parquet/Cassandra

我刚好做过类似的场景,给你一步步拆解Java的实现方案:

第一步:确保依赖到位

首先你的pom.xml里需要包含这些核心依赖(版本根据你的Spark/Cassandra版本调整):

<dependencies>
    <!-- Spark Core & SQL -->
    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-core_2.12</artifactId>
        <version>3.3.0</version>
    </dependency>
    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-sql_2.12</artifactId>
        <version>3.3.0</version>
    </dependency>
    <!-- Spark Cassandra Connector -->
    <dependency>
        <groupId>com.datastax.spark</groupId>
        <artifactId>spark-cassandra-connector_2.12</artifactId>
        <version>3.3.0</version>
    </dependency>
    <!-- Avro & Kafka Avro -->
    <dependency>
        <groupId>org.apache.avro</groupId>
        <artifactId>avro</artifactId>
        <version>1.11.0</version>
    </dependency>
    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-avro_2.12</artifactId>
        <version>3.3.0</version>
    </dependency>
</dependencies>

第二步:完整代码实现

核心是把Avro的GenericRecord转换成Spark SQL能识别的Row,同时把Avro Schema映射为Spark的StructType,再生成DataFrame完成写入操作:

import org.apache.avro.Schema;
import org.apache.avro.generic.GenericRecord;
import org.apache.spark.SparkConf;
import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.api.java.function.Function;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.RowFactory;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.types.DataTypes;
import org.apache.spark.sql.types.StructField;
import org.apache.spark.sql.types.StructType;
import org.apache.spark.streaming.Duration;
import org.apache.spark.streaming.api.java.JavaDStream;
import org.apache.spark.streaming.api.java.JavaStreamingContext;
import org.apache.spark.streaming.kafka.KafkaUtils;
import scala.Tuple2;

import java.util.ArrayList;
import java.util.List;
import java.util.Map;

public class KafkaAvroToParquetCassandra {
    public static void main(String[] args) throws InterruptedException {
        // 1. 初始化Spark配置和Streaming上下文
        SparkConf conf = new SparkConf()
                .setAppName("KafkaAvroToParquetCassandra")
                .setMaster("local[*]") // 生产环境请移除该配置
                .set("spark.cassandra.connection.host", "你的Cassandra主机地址")
                .set("spark.cassandra.connection.port", "9042");

        JavaStreamingContext jssc = new JavaStreamingContext(conf, new Duration(5000));
        SparkSession spark = SparkSession.builder().config(conf).getOrCreate();

        // 2. 初始化Kafka流(复用你的示例代码,修正语法问题)
        Map<String, String> kafkaParams = Map.of(
                "metadata.broker.list", "你的Kafka broker地址",
                "schema.registry.url", "你的Schema Registry地址"
        );
        Map<String, Integer> topicMap = Map.of("你的目标topic", 1);

        JavaPairReceiverInputDStream<String, GenericRecord> kafkaStream = KafkaUtils.createStream(
                jssc,
                String.class,
                GenericRecord.class,
                KafkaAvroDecoder.class,
                KafkaAvroDecoder.class,
                kafkaParams,
                topicMap
        );

        JavaDStream<GenericRecord> msgStream = kafkaStream.map(
                new Function<Tuple2<String, GenericRecord>, GenericRecord>() {
                    @Override
                    public GenericRecord call(Tuple2<String, GenericRecord> tuple2) throws Exception {
                        return tuple2._2();
                    }
                }
        );

        // 3. 处理每个批次的RDD,转为DataFrame并写入存储
        msgStream.foreachRDD(rdd -> {
            if (!rdd.isEmpty()) {
                // 提取Avro Schema并转为Spark StructType
                GenericRecord firstRecord = rdd.first();
                Schema avroSchema = firstRecord.getSchema();
                StructType sparkSchema = convertAvroSchemaToSparkSchema(avroSchema);

                // 将GenericRecord转为Spark Row
                JavaRDD<Row> rowRDD = rdd.map(record -> {
                    List<Object> values = new ArrayList<>();
                    for (Schema.Field field : avroSchema.getFields()) {
                        values.add(record.get(field.name()));
                    }
                    return RowFactory.create(values.toArray());
                });

                // 创建DataFrame并验证
                Dataset<Row> df = spark.createDataFrame(rowRDD, sparkSchema);
                df.show();

                // 写入Parquet文件
                df.write()
                        .mode("append") // 可选:overwrite/ignore/errorIfExists
                        .parquet("/你的Parquet输出路径");

                // 写入Cassandra
                df.write()
                        .format("org.apache.spark.sql.cassandra")
                        .option("keyspace", "你的Cassandra keyspace")
                        .option("table", "你的Cassandra表名")
                        .mode("append")
                        .save();
            }
        });

        // 启动Streaming任务
        jssc.start();
        jssc.awaitTermination();
    }

    // 工具方法:Avro Schema转Spark StructType
    private static StructType convertAvroSchemaToSparkSchema(Schema avroSchema) {
        List<StructField> fields = new ArrayList<>();
        for (Schema.Field field : avroSchema.getFields()) {
            Schema.Type avroType = field.schema().getType();
            switch (avroType) {
                case INT:
                    fields.add(DataTypes.createStructField(field.name(), DataTypes.IntegerType, true));
                    break;
                case STRING:
                    fields.add(DataTypes.createStructField(field.name(), DataTypes.StringType, true));
                    break;
                case LONG:
                    fields.add(DataTypes.createStructField(field.name(), DataTypes.LongType, true));
                    break;
                case BOOLEAN:
                    fields.add(DataTypes.createStructField(field.name(), DataTypes.BooleanType, true));
                    break;
                case FLOAT:
                    fields.add(DataTypes.createStructField(field.name(), DataTypes.FloatType, true));
                    break;
                case DOUBLE:
                    fields.add(DataTypes.createStructField(field.name(), DataTypes.DoubleType, true));
                    break;
                // 如需支持嵌套结构、数组等复杂类型,可在此扩展
                default:
                    throw new IllegalArgumentException("暂不支持的Avro类型: " + avroType);
            }
        }
        return DataTypes.createStructType(fields);
    }
}

关键注意事项

  • Schema一致性:代码假设所有Kafka消息的Avro Schema一致,若存在Schema演进,需额外对接Schema Registry获取最新Schema进行转换。
  • Cassandra表结构:Cassandra表的字段名、数据类型必须与DataFrame完全匹配,否则写入会失败。
  • 生产环境优化:移除setMaster("local[*]"),配置合理的Spark资源,Parquet输出路径可改为HDFS或云存储,Cassandra连接参数建议通过外部配置文件传入。
  • 复杂类型支持:如果你的Avro Schema包含嵌套结构、数组等,需要扩展convertAvroSchemaToSparkSchema方法来适配这些类型。

内容的提问来源于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 09:42:21