【发布时间】: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