PySpark における耐障害性のある分散データセット

PySpark 入門

Benjamin Schmidt

Data Engineer

PySpark における並列化とは何ですか?

  • クラスター内の複数ノードにデータと計算を自動的に並列分散する
  • 複数ノードにわたる大規模データセットの分散処理
  • ワーカーノードはデータを並列処理し、タスクの最後に結合します
  • 大規模処理の高速化(ギガバイト、あるいはテラバイトを想定)

並列化

PySpark 入門

RDD の理解

RDDs または レジリエント分散データセット:

  • ノード障害からの自動復旧を備えたクラスタ全体への分散データコレクション
  • 大規模データに適している
  • map()filter()のような操作で不変に変換でき、collect()paralelize() のようなアクションで結果を取得したりRDDを作成したりできます
PySpark 入門

RDD の作成

# Initialize a Spark session
from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("RDDExample").getOrCreate()

# Create a DataFrame from a csv census_df = spark.read.csv("/census.csv")
# Convert DataFrame to RDD census_rdd = census_df.rdd
# Show the RDD's contents using collect() census_rdd.collect()
PySpark 入門

Collect の表示

# Collect the entire DataFrame into a local Python list of Row objects
data_collected = df.collect()

# Print the collected data
for row in data_collected:
    print(row)
```    
PySpark 入門

RDDとDataFrameの比較

DataFrames

  • 高レベル: 使いやすさを重視して最適化された
  • SQLのLIKE操作: SQLライクなクエリを扱い、より少ないコードで複雑な操作を実行する
  • スキーマ情報: SQLテーブルのような列と型を含む

RDDS

  • 低レベル: より柔軟だが、複雑な操作にはより多くのコード行が必要になる
  • 型安全性: データ型を保持しますが、DataFrame の最適化によるメリットはありません
  • スキーマなし: SQL やリレーショナルデータのような構造化データを扱うのが難しい
  • 大規模な拡張
  • DataFrameと比べて非常に冗長で、分析機能は貧弱
PySpark 入門

便利な関数とメソッド

  • map(): メソッドは、ラムダ関数のような自分で書く関数も含めて、データセット全体に関数を適用します。 rdd.map(map_function)
  • collect(): クラスター全体から次のようなデータを収集します: たとえば rdd.collect()
PySpark 入門

練習しましょう!

PySpark 入門

Preparing Video For Download...