【发布时间】:2018-08-11 10:00:03
【问题描述】:
我正在尝试使用 spark 处理一个大型 cassandra 表(约 4.02 亿个条目和 84 列),但我得到的结果不一致。最初的要求是将一些列从这个表复制到另一个表。复制数据后,我注意到新表中的某些条目丢失了。为了验证我是否计算了大型源表,但我每次都得到不同的值。我在一个较小的表(约 700 万条记录)上尝试了查询,结果很好。
最初,我尝试使用 pyspark 进行计数。这是我的 pyspark 脚本:
spark = SparkSession.builder.appName("Datacopy App").getOrCreate()
df = spark.read.format("org.apache.spark.sql.cassandra").options(table=sourcetable, keyspace=sourcekeyspace).load().cache()
df.createOrReplaceTempView("data")
query = ("select count(1) from data " )
vgDF = spark.sql(query)
vgDF.show(10)
Spark提交命令如下:
~/spark-2.1.0-bin-hadoop2.7/bin/spark-submit --master spark://10.128.0.18:7077 --packages datastax:spark-cassandra-connector:2.0.1-s_2.11 --conf spark.cassandra.connection.host="10.128.1.1,10.128.1.2,10.128.1.3" --conf "spark.storage.memoryFraction=1" --conf spark.local.dir=/media/db/ --executor-memory 10G --num-executors=6 --executor-cores=2 --total-executor-cores 18 pyspark_script.py
上述 spark 提交过程大约需要 90 分钟才能完成。我运行了 3 次,结果如下:
- Spark 迭代 1:402273852
- Spark 迭代 2:402273884
- Spark 迭代 3:402274209
Spark 在整个过程中不显示任何错误或异常。我在 cqlsh 中运行了三次相同的查询并再次得到不同的结果:
- Cqlsh迭代1:402273598
- Cqlsh 迭代 2:402273499
- Cqlsh 迭代 3:402273515
我无法找出为什么我会从同一个查询中得到不同的结果。 Cassandra 系统日志 (/var/log/cassandra/system.log) 仅显示一次以下错误消息:
ERROR [SSTableBatchOpen:3] 2018-02-27 09:48:23,592 CassandraDaemon.java:226 - Exception in thread Thread[SSTableBatchOpen:3,5,main]
java.lang.AssertionError: Stats component is missing for sstable /media/db/datakeyspace/sensordata1-acfa7880acba11e782fd9bf3ae460699/mc-58617-big
at org.apache.cassandra.io.sstable.format.SSTableReader.open(SSTableReader.java:460) ~[apache-cassandra-3.9.jar:3.9]
at org.apache.cassandra.io.sstable.format.SSTableReader.open(SSTableReader.java:375) ~[apache-cassandra-3.9.jar:3.9]
at org.apache.cassandra.io.sstable.format.SSTableReader$4.run(SSTableReader.java:536) ~[apache-cassandra-3.9.jar:3.9]
at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) ~[na:1.8.0_131]
at java.util.concurrent.FutureTask.run(FutureTask.java:266) ~[na:1.8.0_131]
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142) ~[na:1.8.0_131]
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617) [na:1.8.0_131]
at java.lang.Thread.run(Thread.java:748) [na:1.8.0_131]
版本:
- 卡桑德拉 3.9。
- Spark 2.1.0。
- Datastax 的 spark-cassandra-connector 2.0.1
- Scala 2.11 版
集群:
- 具有 3 个工作节点和 1 个主节点的 Spark 设置。
- 3 个工作节点也安装了 cassandra 集群。
- 每个工作节点都有 8 个 CPU 内核和 40 GB RAM。
任何帮助将不胜感激。
【问题讨论】:
-
这不是一个实时数据库(被其他应用程序/进程读取或写入)吗?
-
不,目前只有我在使用数据库。
-
您在写入初始数据时遇到了麻烦吗?我所说的麻烦是指一些网络问题。或者甚至一个 cassandra 没有倒下?你后来修过你的桌子吗?你能告诉我们更多关于表定义的信息吗(例如复制因子和/或压缩策略)。
-
或者它可能只是一个愚蠢的东西,比如过期的 TTL?
-
感谢您的回复。 Cassandra 部署在 gcloud 虚拟机上,因此网络没有问题。该问题与 cassandra 的读取一致性级别有关。
标签: apache-spark cassandra pyspark spark-cassandra-connector