使用scala语言,实现以下功能:消费kafka数据,写入doris数据库
·
以下是一个使用Scala编写的示例代码,实现从Kafka消费数据并写入Doris数据库的功能。代码依赖spark-sql和spark-kafka库,采用Spark Structured Streaming实现流式处理。
依赖配置
在build.sbt中添加以下依赖:
libraryDependencies ++= Seq(
"org.apache.spark" %% "spark-sql" % "3.3.0",
"org.apache.spark" %% "spark-sql-kafka-0-10" % "3.3.0",
"mysql" % "mysql-connector-java" % "8.0.28"
)
Kafka消费与Doris写入代码
import org.apache.spark.sql.{SparkSession, DataFrame}
import org.apache.spark.sql.functions._
import org.apache.spark.sql.streaming.Trigger
object KafkaToDoris {
def main(args: Array[String]): Unit = {
val spark = SparkSession.builder()
.appName("KafkaToDoris")
.master("local[*]") // 生产环境需替换为集群地址
.getOrCreate()
// Kafka配置
val kafkaParams = Map(
"kafka.bootstrap.servers" -> "kafka-server:9092",
"subscribe" -> "your_topic",
"startingOffsets" -> "latest"
)
// 从Kafka读取数据
val kafkaDF = spark
.readStream
.format("kafka")
.options(kafkaParams)
.load()
// 解析Kafka消息(假设为JSON格式)
val parsedDF = kafkaDF
.select(from_json(col("value").cast("string"), schema).alias("data"))
.select("data.*")
// 定义Doris JDBC连接参数
val dorisUrl = "jdbc:mysql://doris-fe:9030/database_name"
val dorisProps = new java.util.Properties()
dorisProps.put("user", "username")
dorisProps.put("password", "password")
dorisProps.put("driver", "com.mysql.jdbc.Driver")
// 写入Doris
val query = parsedDF.writeStream
.foreachBatch { (batchDF: DataFrame, batchId: Long) =>
batchDF.write
.mode("append")
.jdbc(dorisUrl, "target_table", dorisProps)
}
.trigger(Trigger.ProcessingTime("30 seconds"))
.start()
query.awaitTermination()
}
}
关键配置说明
- Kafka参数:需根据实际环境调整
bootstrap.servers和topic名称。 - 数据解析:示例假设Kafka消息为JSON格式,需根据实际数据格式调整解析逻辑。
- Doris连接:Doris兼容MySQL协议,直接使用JDBC连接。需替换
doris-fe:9030为实际的FE节点地址。 - 写入模式:
append表示追加数据,支持overwrite等模式。
性能优化建议
- 调整
Trigger.ProcessingTime控制写入频率。 - 在Doris侧合理设置分区分桶,提升写入和查询性能。
- 对于大批量写入,可考虑启用Spark的批量提交功能。
注意事项
- Doris需提前创建好目标表,且表结构需与Spark DataFrame的Schema匹配。
- 生产环境建议配置Checkpoint保证Exactly-Once语义。
更多推荐




所有评论(0)