始める無料で始める

データにSparkをひと振り

前の演習では、Sparkから pandas へデータを移す方法を見ました。今度はその逆に、pandas のDataFrameをSparkクラスターに取り込んでみましょう。SparkSession クラスにはそのためのメソッドがあります。

.createDataFrame() メソッドは、pandas のDataFrameを受け取り、SparkのDataFrameを返します。

このメソッドの出力はローカルに保持され、SparkSession のカタログには登録されません。つまり、Spark DataFrameの各種メソッドは使えますが、ほかのコンテキストからはデータにアクセスできません。

たとえば、そのDataFrameを参照するSQLクエリ(.sql() メソッド使用)はエラーになります。この方法でアクセスしたい場合は、まず 一時テーブル として保存する必要があります。

これには、Spark DataFrameの .createTempView() メソッドを使います。引数は登録したい一時テーブル名の1つだけです。このメソッドはDataFrameをカタログ内のテーブルとして登録しますが、一時テーブルなので、そのSpark DataFrameを作成した特定の SparkSession からしかアクセスできません。

.createOrReplaceTempView() というメソッドもあります。これは、未登録なら新しい一時テーブルを作成し、既に存在する場合は安全に置き換えます。重複テーブルによる問題を避けるため、このメソッドを使います。

図を見て、Sparkのデータ構造同士がどのように相互作用するかを確認しましょう。

作業環境にはすでに spark という SparkSession が用意され、numpynppandaspd としてインポート済みです。

この演習はコースの一部です

PySpark入門

コースを見る

演習の手順

  • 乱数で pandas のDataFrameを作成するコードはすでに用意され、pd_temp として保存されています。
  • Sparkの .createDataFrame()pd_temp を引数にして呼び出し、spark_temp というSpark DataFrameを作成してください。
  • Sparkクラスター内のテーブル一覧を確認し、新しいDataFrameが存在しないことを確かめてください。spark.catalog.listTables() が使えます。
  • 直前に作成した spark_temp.createOrReplaceTempView() で一時テーブルとして登録してください。一時テーブル名は "temp" にします。テーブル名はこのメソッドの唯一の引数として指定する点に注意してください!
  • もう一度、テーブル一覧を確認しましょう。

実践的なインタラクティブ演習

このサンプルコードを完成させて、この演習に挑戦してみましょう。

# Create pd_temp
pd_temp = pd.DataFrame(np.random.random(10))

# Create spark_temp from pd_temp
spark_temp = ____

# Examine the tables in the catalog
print(____)

# Add spark_temp to the catalog
spark_temp.____

# Examine the tables in the catalog again
print(____)
コードを編集して実行