Scala实现Kafka流JSON数据写入Cassandra遇阻求助
解决Kafka流JSON数据写入Cassandra的问题
嘿,我来帮你搞定这个问题!看了你的代码,能发现几个核心的小问题,咱们一步步来修正,最后给你一套能跑通的完整方案:
先说说你现有代码里的坑
- JSON处理完全错了:你用
flatMap(_.split(","))来拆JSON字符串可不行啊!JSON是有结构的,不是简单用逗号分割的,得用专门的JSON解析库来提取vehicleid和vehicletype字段。 - Case Class命名不符合Scala规范:Scala里类名首字母得大写,
event改成Event才对,不然后续可能会有隐性问题。 - 缺少Cassandra连接配置:
saveToCassandra可不是直接就能用的,得先配置Cassandra的连接信息,项目也得引入对应的依赖包。
完整实现步骤
第一步:先把依赖配好
在你的build.sbt里加上这些依赖(根据你的Spark版本调整,我这里用的是Spark 3.3.x的版本):
libraryDependencies ++= Seq( "org.apache.spark" %% "spark-streaming" % "3.3.0" % Provided, "org.apache.spark" %% "spark-streaming-kafka-0-10" % "3.3.0", "com.datastax.spark" %% "spark-cassandra-connector" % "3.3.0", "org.json4s" %% "json4s-native" % "4.0.6" // 用这个JSON解析库比较方便,也可以换Jackson )
第二步:修正后的完整代码
下面是能直接运行的代码,包含了JSON解析、Cassandra连接和写入逻辑:
import org.apache.spark.SparkConf import org.apache.spark.streaming.{Seconds, StreamingContext} import org.apache.spark.streaming.kafka010.{ConsumerStrategies, KafkaUtils, LocationStrategies} import com.datastax.spark.connector.streaming._ import org.json4s._ import org.json4s.native.JsonMethods._ // 要和Cassandra表结构对应,类名首字母大写 case class Event(vehicleid: String, vehicletype: String) object KafkaToCassandra { def main(args: Array[String]): Unit = { val conf = new SparkConf() .setAppName("KafkaToCassandra") .setMaster("local[*]") // 生产环境一定要删掉这句,用集群配置 // 配置Cassandra连接信息 .set("spark.cassandra.connection.host", "你的Cassandra地址") .set("spark.cassandra.connection.port", "9042") // 如果Cassandra开了认证,加上下面两行 // .set("spark.cassandra.auth.username", "你的用户名") // .set("spark.cassandra.auth.password", "你的密码") // 每5秒处理一个批次,根据你的需求调整 val ssc = new StreamingContext(conf, Seconds(5)) // Kafka的配置参数 val kafkaParams = Map[String, Object]( "bootstrap.servers" -> "你的Kafka地址:9092", "key.deserializer" -> "org.apache.kafka.common.serialization.StringDeserializer", "value.deserializer" -> "org.apache.kafka.common.serialization.StringDeserializer", "group.id" -> "kafka-to-cassandra-group", "auto.offset.reset" -> "latest", "enable.auto.commit" -> (false: java.lang.Boolean) ) // 你要消费的Kafka主题 val topics = Array("你的Kafka主题名称") // 用新版的Kafka 0-10 API创建流,旧的createDirectStream已经过时啦 val kafkaStream = KafkaUtils.createDirectStream[String, String]( ssc, LocationStrategies.PreferConsistent, ConsumerStrategies.Subscribe[String, String](topics, kafkaParams) ) // 解析Kafka里的JSON数据 val events = kafkaStream.map(record => { val jsonStr = record.value() // 用json4s解析JSON,自动映射到Event类 implicit val formats = DefaultFormats parse(jsonStr).extract[Event] }) // 打印当前批次的数据,方便测试(可选) events.foreachRDD(rdd => { println("=== 当前批次的事件数据 ===") rdd.foreach(println) }) // 写入Cassandra,参数是keyspace名称和表名称 events.saveToCassandra("你的Cassandra键空间", "你的表名") // 启动流处理程序,等待终止 ssc.start() ssc.awaitTermination() } }
第三步:准备Cassandra表
得先在Cassandra里建好对应的表,结构要和Event类完全匹配:
-- 创建键空间(如果没有的话) CREATE KEYSPACE IF NOT EXISTS your_keyspace WITH replication = {'class': 'SimpleStrategy', 'replication_factor': 1}; -- 切换到这个键空间 USE your_keyspace; -- 创建表,vehicleid作为主键 CREATE TABLE IF NOT EXISTS vehicle_events ( vehicleid TEXT PRIMARY KEY, vehicletype TEXT );
几个要注意的点
- JSON格式要对:Kafka里发送的JSON必须是
{"vehicleid":"v1", "vehicletype":"car"}这种结构,不然解析会报错。 - 版本要兼容:Spark、Spark Cassandra Connector、Kafka的版本必须对应,比如Spark 3.3.x就得用Connector 3.3.x,别乱搭配。
- 生产环境调整:把
setMaster("local[*]")删掉,配置好Spark集群的资源,Cassandra的连接参数也要改成生产环境的地址。
内容的提问来源于stack exchange,提问作者jAi
相关产品推荐
相关产品推荐

