首页
学习
活动
专区
圈层
工具
发布
社区首页 >问答首页 >数据源io.pivotal.greenplum.spark.GreenplumRelationProvider不支持流写入。

数据源io.pivotal.greenplum.spark.GreenplumRelationProvider不支持流写入。
EN

Stack Overflow用户
提问于 2019-04-04 13:51:35
回答 2查看 392关注 0票数 0

我试图读取卡夫卡的数据,并将其上传到格林梅利的数据库使用火花。我使用的是格林梅-火花连接器,但我正在获取数据源,io.pivotal.greenplum.spark.GreenplumRelationProvider不支持流写入。是否格林梅源不支持流媒体数据?我可以在网站上看到“连续ETL管道(流)”。

我曾尝试将数据源命名为“绿梅”,并将"io.pivotal.greenplum.spark.GreenplumRelationProvider“转换为.format(”数据源“)

代码语言:javascript
复制
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()
EN

回答 2

Stack Overflow用户

回答已采纳

发布于 2019-04-04 18:05:52

格林梅星火结构化流

演示如何使用JDBC将writeStream API与GPDB一起使用。

下面的代码块使用速率流源进行读取,并使用基于JDBC的接收器批量流到GPDB

基于批量的流

代码语言:javascript
复制
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

代码语言:javascript
复制
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
票数 0
EN

Stack Overflow用户

发布于 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

票数 1
EN
页面原文内容由Stack Overflow提供。腾讯云小微IT领域专用引擎提供翻译支持
原文链接:

https://stackoverflow.com/questions/55517663

复制
相关文章

相似问题

领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档