【问题标题】:Spark-Cassandra , How to get data based on QuerySpark-Cassandra , 如何基于 Query 获取数据
【发布时间】:2021-09-16 17:56:30
【问题描述】:

我有一个非常大的 Cassandra 表,现在我与以下代码建立了 spark-Cassandra 连接。

import pandas as pd
import numpy as np
from pyspark import *
import os
from pyspark.sql import SQLContext


os.environ['PYSPARK_SUBMIT_ARGS'] = '--packages  com.datastax.spark:spark-cassandra-connector_2.12:3.0.1 --conf spark.cassandra.connection.host=127.0.0.1 pyspark-shell'
conf = SparkConf().set("spark.cassandra.connection.host", "127.0.0.1").set("spark.cassandra.connection.port", "9042").setAppName("Sentinel").setMaster("spark://Local:7077")
sc = SparkContext(conf=conf)
sqlContext = SQLContext(sc)

table_df = sqlContext.read\
        .format("org.apache.spark.sql.cassandra")\
        .options(table='movies', keyspace='movie_lens')\
        .load()\
        

主键是 Movie_id,它是一个整数。 .load() 将整个表加载到内存中,这是我想避免的。我得到的一种方法是使用过滤器

table_df = sqlContext.read\
        .format("org.apache.spark.sql.cassandra")\
        .options(table='movies', keyspace='movie_lens')\
        .load()\
        .filter("movie_id = 37032")

但是过滤器实际上会阻止将整个表加载到内存中吗?还是先加载然后过滤。 此外,我必须查询许多 ID。假设我需要 1000 个 ID,并且每天 ID 都在不断变化。那怎么办呢?

【问题讨论】:

    标签: apache-spark pyspark cassandra spark-cassandra-connector


    【解决方案1】:

    是的,如果您对分区键进行查询,Spark Cassandra 连接器将执行所谓的“谓词下推”,并且只会从特定查询中加载数据(.load 函数只会加载元数据,实际当您确实需要数据来执行操作时,数据加载将第一次发生)。关于何时在 Spark Cassandra 连接器中发生谓词下推,有 well documented 规则。您也可以通过运行table_df.explain() 来检查这一点,并在PushedFilters 部分查找标有星号* 的过滤器。

    如果您需要查找多个 ID,则可以使用 .isin 过滤器,但实际上不建议使用 Cassandra。最好创建一个带有 ID 的数据帧,并使用 Cassandra 数据帧执行所谓的Direct Join(它从 SCC 2.5 开始适用于数据帧,或更早版本适用于 RDD)。我有一个 lengthy blog post 在加入 Cassandra 中的数据

    【讨论】:

    • 您好,非常感谢您的快速回复。我只是有一个后续问题。你能告诉我 .filter() 和 .where() 有什么区别吗?据我所知,他们做同样的事情。
    • 它们之间没有区别,只是功能相同——出于历史原因
    猜你喜欢
    • 2015-09-13
    • 1970-01-01
    • 2016-02-26
    • 2017-06-15
    • 1970-01-01
    • 1970-01-01
    • 2017-04-17
    • 2020-02-21
    • 2016-08-14
    相关资源
    最近更新 更多