【发布时间】:2017-11-07 13:42:10
【问题描述】:
我正在尝试获取 pyspark 中 cassandra 表的分区键的不同值。但是,pyspark 似乎不理解我,并且完全迭代所有数据(很多)而不是查询索引。
这是我使用的代码,对我来说看起来很简单:
from pyspark.sql import SparkSession
spark = SparkSession \
.builder \
.appName("Spark! This town not big enough for the two of us.") \
.getOrCreate()
ct = spark.read\
.format("org.apache.spark.sql.cassandra")\
.options(table="avt_sensor_data", keyspace="ipe_smart_meter")\
.load()
all_sensors = ct.select("machine_name", "sensor_name")\
.distinct() \
.collect()
“machine_name”和“sensor_name”列共同构成分区键(完整架构见下文)。在我看来,这应该是超快的,事实上,如果我在 cql 中执行这个查询只需要几秒钟:
select distinct machine_name,sensor_name from ipe_smart_meter.avt_sensor_data;
但是,火花作业大约需要 10 个小时才能完成。从 spark 告诉我的计划来看,它似乎真的想要迭代所有数据:
== Physical Plan ==
*HashAggregate(keys=[machine_name#0, sensor_name#1], functions=[], output=[machine_name#0, sensor_name#1])
+- Exchange hashpartitioning(machine_name#0, sensor_name#1, 200)
+- *HashAggregate(keys=[machine_name#0, sensor_name#1], functions=[], output=[machine_name#0, sensor_name#1])
+- *Scan org.apache.spark.sql.cassandra.CassandraSourceRelation@2ee2f21d [machine_name#0,sensor_name#1] ReadSchema: struct<machine_name:string,sensor_name:string>
我不是专家,但在我看来这不像是“使用 cassandra 索引”。
我做错了什么?有没有办法告诉 spark 委派从 cassandra 获取不同值的任务?任何帮助将不胜感激!
如果有帮助,这里是底层 cassandra 表的架构描述:
CREATE KEYSPACE ipe_smart_meter WITH replication = {'class': 'SimpleStrategy', 'replication_factor': '2'} AND durable_writes = true;
CREATE TABLE ipe_smart_meter.avt_sensor_data (
machine_name text,
sensor_name text,
ts timestamp,
id bigint,
value double,
PRIMARY KEY ((machine_name, sensor_name), ts)
) WITH CLUSTERING ORDER BY (ts DESC)
AND bloom_filter_fp_chance = 0.01
AND caching = {'keys': 'ALL', 'rows_per_partition': 'NONE'}
AND comment = '[PRODUCTION] Table for raw data from AVT smart meters.'
AND compaction = {'class': 'org.apache.cassandra.db.compaction.DateTieredCompactionStrategy', 'max_threshold': '32', 'min_threshold': '4'}
AND compression = {'chunk_length_in_kb': '64', 'class': 'org.apache.cassandra.io.compress.LZ4Compressor'}
AND crc_check_chance = 1.0
AND dclocal_read_repair_chance = 0.1
AND default_time_to_live = 0
AND gc_grace_seconds = 864000
AND max_index_interval = 2048
AND memtable_flush_period_in_ms = 0
AND min_index_interval = 128
AND read_repair_chance = 0.0
AND speculative_retry = '99PERCENTILE';
【问题讨论】:
标签: apache-spark cassandra pyspark spark-cassandra-connector