以下是一个使用Scala编写的示例代码,实现从Kafka消费数据并写入Doris数据库的功能。代码依赖spark-sqlspark-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.serverstopic名称。
  • 数据解析:示例假设Kafka消息为JSON格式,需根据实际数据格式调整解析逻辑。
  • Doris连接:Doris兼容MySQL协议,直接使用JDBC连接。需替换doris-fe:9030为实际的FE节点地址。
  • 写入模式append表示追加数据,支持overwrite等模式。

性能优化建议

  • 调整Trigger.ProcessingTime控制写入频率。
  • 在Doris侧合理设置分区分桶,提升写入和查询性能。
  • 对于大批量写入,可考虑启用Spark的批量提交功能。

注意事项

  • Doris需提前创建好目标表,且表结构需与Spark DataFrame的Schema匹配。
  • 生产环境建议配置Checkpoint保证Exactly-Once语义。
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐