【问题标题】:PySpark + Cassandra: Getting distinct values of partition keyPySpark + Cassandra:获取分区键的不同值
【发布时间】: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


    【解决方案1】:

    似乎自动 cassandra 服务器端下推谓词仅在选择、过滤或排序时有效。

    https://github.com/datastax/spark-cassandra-connector/blob/master/doc/14_data_frames.md

    所以,如果是您的 distinct(),spark 会获取所有行,然后是 distinct()。

    解决方案 1

    你说你的 cql select distinct... 已经超快了。我猜分区键的数量相对较少(machine_name 和 sensor_name 的组合)和很多 'ts'。

    所以,最简单的解决方案就是使用 cql(例如,cassandra-driver)。

    解决方案 2

    由于 cassandra 是一个查询优先的数据库,因此只需再创建一个表,该表仅包含不同查询所需的分区键。

    CREATE TABLE ipe_smart_meter.avt_sensor_name_machine_name (
        machine_name text,
        sensor_name text,
        PRIMARY KEY ((machine_name, sensor_name))
    );
    

    然后,每次在原始表中插入一行时,将 machine_name 和 sensor_name 插入到新表中。 由于它只有分区键,因此这是您查询的自然不同表。只需获取所有行。也许超快。无需区分流程。

    解决方案 3

    我认为解决方案 2 是最好的。但是,如果您不想对一条记录进行两次插入,另一种解决方案是更改您的表并创建一个物化视图表。

    CREATE TABLE ipe_smart_meter.ipe_smart_meter.avt_sensor_data (
        machine_name text,
        sensor_name text,
        ts timestamp,
        id bigint,
        value double,
        dist_hint_num smallint,
        PRIMARY KEY ((machine_name, sensor_name), ts)
    ) WITH CLUSTERING ORDER BY (ts DESC)
    ;
    
    CREATE MATERIALIZED VIEW IF NOT EXISTS ipe_smart_meter.avt_sensor_data_mv AS
      SELECT
        machine_name
        ,sensor_name
        ,ts
        ,dist_hint_num
      FROM ipe_smart_meter.avt_sensor_data
      WHERE
        machine_name IS NOT NULL
        AND sensor_name IS NOT NULL
        AND ts IS NOT NULL
        AND dist_hint_num IS NOT NULL
      PRIMARY KEY ((dist_hint_num), machine_name, sensor_name, ts)
      WITH
      AND CLUSTERING ORDER BY (machine_name ASC, sensor_name DESC, ts DESC)
    ;
    

    dist_hint_num 列用于限制您的查询迭代和分发记录的分区总数。

    例如,从 0 到 15。随机整数 random.randint(0, 15) 或基于哈希的整数 hash_func(machine_name + sensor_name) % 16 都可以。 然后,当您查询如下。 cassandra 仅从 16 个分区获取所有记录,这可能比您当前的情况更有效。

    但是,无论如何,必须先读取所有记录,然后再读取distinct()(随机播放)。不节省空间。我认为这不是一个好的解决方案。

    functools.reduce(
        lambda df, dist_hint_num: df.union(
            other=spark_session.read.format(
                'org.apache.spark.sql.cassandra',
            ).options(
                keyspace='ipe_smart_meter',
                table='avt_sensor_data_mv',
            ).load().filter(
                col('dist_hint_num') == expr(
                    f'CAST({dist_hint_num} AS SMALLINT)'
                )
            ).select(
                col('machine_name'),
                col('sensor_name'),
            ),
        ),
        range(0, 16),
        spark_session.createDataFrame(
            data=(),
            schema=StructType(
                fields=(
                    StructField(
                        name='machine_name',
                        dataType=StringType(),
                        nullable=False,
                    ),
                    StructField(
                        name='sensor_name',
                        dataType=StringType(),
                        nullable=False,
                    ),
                ),
            ),
        ),
    ).distinct().persist().alias(
        'df_all_machine_sensor',
    )
    

    【讨论】:

    • 谢谢,实际上我最终使用了解决方案 1,因为我没有对该数据库的写入权限。
    猜你喜欢
    • 2021-12-12
    • 1970-01-01
    • 1970-01-01
    • 2020-10-05
    • 1970-01-01
    • 2017-03-15
    • 2020-02-16
    • 2016-03-31
    • 2015-06-21
    相关资源
    最近更新 更多