Spark读写dorishivekafka
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
