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
相关产品推荐
相关产品推荐

