我试图读取卡夫卡的数据,并将其上传到格林梅利的数据库使用火花。我使用的是格林梅-火花连接器,但我正在获取数据源,io.pivotal.greenplum.spark.GreenplumRelationProvider不支持流写入。是否格林梅源不支持流媒体数据?我可以在网站上看到“连续ETL管道(流)”。
我曾尝试将数据源命名为“绿梅”,并将"io.pivotal.greenplum.spark.GreenplumRelationProvider“转换为.format(”数据源“)
val EventStream = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", args(0))
.option("subscribe", args(1))
.option("startingOffsets", "earliest")
.option("failOnDataLoss", "false")
.load
val gscWriteOptionMap = Map(
"url" -> "link for greenplum",
"user" -> "****",
"password" -> "****",
"dbschema" -> "dbname"
)
val stateEventDS = EventStream
.selectExpr("CAST(key AS String)", "*****(value)")
.as[(String, ******)]
.map(_._2)
val EventOutputStream = stateEventDS.writeStream
.format("io.pivotal.greenplum.spark.GreenplumRelationProvider")
.options(gscWriteOptionMap)
.start()
assetEventOutputStream.awaitTermination()发布于 2019-04-04 18:05:52
格林梅星火结构化流
演示如何使用JDBC将writeStream API与GPDB一起使用。
下面的代码块使用速率流源进行读取,并使用基于JDBC的接收器批量流到GPDB
基于批量的流
import org.apache.spark.sql.streaming._
import org.apache.spark.sql.streaming.Trigger
import org.apache.spark.streaming.Seconds
import org.apache.spark.streaming.StreamingContext
import scala.concurrent.duration._
val sq = spark.
readStream.
format("rate").
load.
writeStream.
format("myjdbc").
option("checkpointLocation", "/tmp/jdbc-checkpoint").
trigger(Trigger.ProcessingTime(10.seconds)).
start基于记录的流
这使用了ForeachWriter
import org.apache.spark.sql.streaming._
import org.apache.spark.sql.streaming.Trigger
import org.apache.spark.streaming.Seconds
import org.apache.spark.streaming.StreamingContext
import scala.concurrent.duration._
val url="jdbc:postgresql://gsc-dev:5432/gpadmin"
val user ="gpadmin"
val pwd = "changeme"
val jdbcWriter = new JDBCSink(url,user, pwd)
val sq = spark.
readStream.
format("rate").
load.
writeStream.
format(jdbcWriter).
option("checkpointLocation", "/tmp/jdbc-checkpoint").
trigger(Trigger.ProcessingTime(10.seconds)).
start发布于 2019-04-04 15:46:27
你使用的是什么版本的GPDB /火花?你可以绕过火花,以支持格林梅-卡夫卡连接器。
https://gpdb.docs.pivotal.io/5170/greenplum-kafka/overview.html
在早期版本中,格林梅-斯帕克连接器公开了一个名为io.pivotal.greenplum.spark.GreenplumRelationProvider的火花数据源,将数据从格林梅利数据库读取到一个DataFrame中。
在以后的版本中,连接器公开了一个名为greenplum的Spark数据源,用于在Spark和之间传输数据。
应该是--
val EventOutputStream =stateEventDS.write.format(“绿梅”) .options(gscWriteOptionMap) .save()
请参阅:gpdb.html
https://stackoverflow.com/questions/55517663
复制相似问题