改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

Spark读写dorishivekafka

ckckck2025年1月10日1 浏览

Spark 与 Doris、Hive、Kafka 集成示例

1. pom.xml 依赖配置

xml 复制代码
<dependency>
    <groupId>org.apache.spark</groupId>
    <artifactId>spark-core_2.12</artifactId>
    <version>3.2.1</version>
</dependency>
<dependency>
    <groupId>org.apache.spark</groupId>
    <artifactId>spark-sql_2.12</artifactId>
    <version>3.2.1</version>
    <scope>provided</scope>
</dependency>
<dependency>
    <groupId>org.apache.doris</groupId>
    <artifactId>spark-doris-connector-3.2_2.12</artifactId>
    <version>1.3.0</version>
</dependency>
<dependency>
    <groupId>org.apache.spark</groupId>
    <artifactId>spark-hive_2.12</artifactId>
    <version>3.2.1</version>
    <!-- <scope>provided</scope> -->
</dependency>
<dependency>
    <groupId>org.apache.spark</groupId>
    <artifactId>spark-sql-kafka-0-10_2.12</artifactId>
    <version>3.2.1</version>
</dependency>
<dependency>
    <groupId>org.apache.spark</groupId>
    <artifactId>spark-streaming_2.12</artifactId>
    <version>3.2.1</version>
    <scope>provided</scope>
</dependency>
<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>3.2.3</version>
</dependency>
<dependency>
    <groupId>org.apache.spark</groupId>
    <artifactId>spark-streaming-kafka-0-10_2.12</artifactId>
    <version>3.2.1</version>
</dependency>

2. SparkSQL 读取 Doris 数据写入 Hive

scala 复制代码
val spark = SparkSession.builder()
  .appName("Doris Reader")
  .config("spark.sql.warehouse.dir", "hdfs://xxxx/user/hive/warehouse")
  .config("hive.metastore.uris", "thrift://xxxx:9083")
  .config("hive.exec.dynamic.partition", true)
  .config("hive.exec.dynamic.partition.mode", "nonstrict")
  .config(conf)
  .enableHiveSupport()
  .getOrCreate()

val syncDate: String = args(0)

// 读取 Doris 分区数据
val df = spark
  .read
  .format("doris")
  .option("doris.table.identifier", "xxx.xxx")
  .option("doris.fenodes", "xxx:xxx")
  .option("user", "....")
  .option("password", "****")
  .option("doris.filter.query", s"pDate = '$syncDate'")
  .load()

// 方式一:df 数据写入 Hive
df.write.format("hive").mode(SaveMode.Overwrite).partitionBy("pdate").saveAsTable("mdz_doris.dwd_vehicleDriverStatusInfoTest")

// 方式二:df 写入 Hive
df.drop("pDate").createOrReplaceTempView("result")
spark.sql(s"insert overwrite mdz_doris.dwd_vehicledriverstatusinfotest PARTITION(pdate='$syncDate') select * from result")

3. SparkSQL 读取 Hive 数据写入 Doris

scala 复制代码
val spark = SparkSession.builder()
  .appName("Spark Hive to Doris")
  .config("spark.sql.warehouse.dir", "hdfs://xxxxx/user/hive/warehouse")
  .config("hive.metastore.uris", "thrift://xxxxx:9083")
  .config(conf)
  .enableHiveSupport()
  .getOrCreate()

val syncDate = args(0)

// 读取 Hive 分区表数据
val hivePartitionTableDF = spark.sql(s"SELECT * FROM mdz_doris.dwd_vehicledriverstatusinfotest WHERE pdate ='$syncDate'")

// 将数据写入 Doris
hivePartitionTableDF
  .write
  .format("doris")
  .option("doris.table.identifier", "xxxx.xxxx")
  .option("doris.fenodes", "xxxx:xxx")
  .option("user", "xxxx")
  .option("password", "****")
  .option("doris.write.fields", "column1,column2,column3,....")
  .save()

spark.stop()

4. SparkSQL 读取 Hive 数据写入 Kafka

scala 复制代码
val spark = SparkSession.builder()
  .appName("HiveToKafka")
  .config("spark.sql.warehouse.dir", "hdfs://xxxx/user/hive/warehouse")
  .config("hive.metastore.uris", "thrift://xxxx:9083")
  .config(conf)
  .enableHiveSupport()
  .getOrCreate()

import spark.implicits._

val syncDate = args(0)
val kafkaBrokers = "***********"

val df1 = spark.sql(
    s"""
      |select * from
      |database.table
      |where pdate= '$syncDate'
      |""".stripMargin)
  .toJSON
  .toDF("value")

// 定义 Kafka 配置
val props = new Properties()
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaBrokers) // Kafka 集群地址
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer")
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer")

// 将 DataFrame 发送到 Kafka
df1.foreachPartition { (partition: Iterator[Row]) =>
  val producer = new KafkaProducer[String, String](props)
  partition.foreach { row =>
    val value = row.getAs[String]("value")
    val record = new ProducerRecord[String, String]("target_topic", value)
    producer.send(record, new Callback {
      override def onCompletion(metadata: RecordMetadata, exception: Exception): Unit = {
        if (exception == null) {
          System.out.println("partition: " + metadata.partition() + " offset: " + metadata.offset())
        } else {
          exception.printStackTrace()
        }
      }
    })
  }
  producer.close()
}

5. Spark 读取 Kafka 数据写入 Hive(结构化流)

scala 复制代码
val kafkaDF = spark
  .readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "xxxx")
  .option("subscribe", topicName)
  .option("failOnDataLoss", false) // 如果读取数据源时,发现数据突然缺失,比如被删,则是否马上抛出异常
  .option("fetchOffset.numRetries", 3) // 获取消息的偏移量时,最多进行的重试次数
  .option("maxOffsetsPerTrigger", 500) // 用于限流,限定每次读取数据的最大条数,不指定则是 as fast as possible
  .option("startingOffsets", "earliest") // 第一次消费时,读取 Kafka 数据的位置
  .load()

import spark.implicits._

val schema: StructType = StructType(
  Array(
    StructField("timestamp", StringType, true),
    StructField("vehicleid", StringType, true),
    // 其他字段
  )
)

val parsedDF = kafkaDF.selectExpr("CAST(value AS STRING)")
  .select(from_json($"value", schema).as("data"))
  .select($"data.*")

val syncDate = args(0)

val query = parsedDF.writeStream
  .queryName("kafka2Hive")
  .outputMode(OutputMode.Append())
  .option("checkpointLocation", "hdfs://xxxx/user/hive/ckp") // 用来保存 offset,用该目录来绑定对应的 offset
  .foreachBatch((df: Dataset[Row], batchId: Long) => {
    myWriteFun(df, batchId, "database.table", syncDate)
  })
  .trigger(Trigger.ProcessingTime(2000))
  .start()

def myWriteFun(df: Dataset[Row], batchId: Long, tableName: String, pDate: String): Unit = {
  println("BatchId" + batchId)
  if (df.count() != 0) {
    df.persist()
    df.write.format("hive").mode(SaveMode.Append).partitionBy(pDate).saveAsTable(tableName)
    df.unpersist()
  }
}

6. Spark 读取 Kafka 写入 Doris(结构化流)

scala 复制代码
val kafkaSource = spark
  .readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "****************")
  .option("subscribe", sourceTopic)
  .option("failOnDataLoss", false) // 如果读取数据源时,发现数据突然缺失,比如被删,则是否马上抛出异常
  .option("startingOffsets", "earliest") // 第一次消费时,读取 Kafka 数据的位置
  .load()

kafkaSource
  .selectExpr("CAST(value as STRING)")
  .writeStream
  .format("doris")
  .option("checkpointLocation", "hdfs://xxxx/user/kafka/ckp")
  .option("doris.table.identifier", "xxxx.xxxx")
  .option("doris.fenodes", "xxxxx:xxxx")
  .option("user", "xxxx")
  .option("password", "****")
  .option("max_filter_ratio", "0.1")
  .option("doris.sink.streaming.passthrough", "true")
  .option("doris.sink.properties.format", "json")
  .option("doris.write.fields", "column1,column2,column2,....")
  .start()
  .awaitTermination()
示例图片

end