【发布时间】:2018-11-11 03:25:43
【问题描述】:
当从一个表读取数据并在 cassandra 中将其写入另一个表时,谁能解释一下 spark 的内部工作原理。
这是我的用例:
我正在通过 kafka 主题将来自 IOT 平台的数据摄取到 cassandra。我有一个小的 python 脚本,它解析来自 kafka 的每条消息以获取它所属的表名,准备一个查询并使用 datastax 的 cassandra-driver for python 将其写入 cassandra。使用该脚本,我可以每分钟将大约 300000 条记录 摄取到 cassandra 中。但是我的传入数据速率是每分钟 510000 条记录,因此 kafka 消费者延迟不断增加。
Python 脚本已经在对 cassandra 进行并发调用。如果我增加 python 执行器的数量,cassandra-driver 开始失败,因为 cassandra 节点对其不可用。我假设我在那里打的每秒 cassandra 调用次数是有限制的。这是我收到的错误消息:
ERROR Operation failed: ('Unable to complete the operation against any hosts', {<Host: 10.128.1.3 datacenter1>: ConnectionException('Pool is shutdown',), <Host: 10.128.1.1 datacenter1>: ConnectionException('Pool is shutdown',)})"
最近,我运行了一个 pyspark 作业,将数据从一个表中的几列复制到另一个。该表中有大约 1.68 亿条记录。 Pyspark 作业在大约 5 小时内完成。因此它每分钟处理超过 550000 条记录。
这是我正在使用的 pyspark 代码:
df = spark.read\
.format("org.apache.spark.sql.cassandra")\
.options(table=sourcetable, keyspace=sourcekeyspace)\
.load().cache()
df.createOrReplaceTempView("data")
query = ("select dev_id,datetime,DATE_FORMAT(datetime,'yyyy-MM-dd') as day, " + field + " as value from data " )
vgDF = spark.sql(query)
vgDF.show(50)
vgDF.write\
.format("org.apache.spark.sql.cassandra")\
.mode('append')\
.options(table=newtable, keyspace=newkeyspace)\
.save()
版本:
- 卡桑德拉 3.9。
- Spark 2.1.0。
- Datastax 的 spark-cassandra-connector 2.0.1
- Scala 2.11 版
集群:
- 具有 3 个工作节点和 1 个主节点的 Spark 设置。
- 3 个工作节点也安装了 cassandra 集群。 (每个 cassandra 节点都有一个 spark 工作节点)
- 允许每个工作人员使用 10 GB 内存和 3 个内核。
所以我想知道:
spark 是否首先从 cassandra 读取所有数据,然后将其写入新表,或者 spark cassandra 连接器中是否有某种优化,允许它在 cassandra 表中移动数据而不读取所有记录?
如果我将我的 python 脚本替换为一个 spark 流作业,在该作业中我解析数据包以获取 cassandra 的表名,这将有助于我更快地将数据摄取到 cassandra 中吗?
【问题讨论】:
标签: apache-spark pyspark cassandra cassandra-3.0 spark-cassandra-connector