データ接続と統合Python (Spark)PySparkリファレンス概念:列

注: 以下の翻訳の正確性は検証されていません。AIPを利用して英語版の原文から機械的に翻訳されたものです。

概念:列

このドキュメントの例に沿って進めるには、from pyspark.sql import functions as F を追加してください。

列は PySpark クラス pyspark.sql.Column によって管理されます。既存の列を直接参照するか、既存の列から式を導出するたびに、列のインスタンスが作成されます。列は、次のいずれかの方法で参照できます。

  • F.col("column_name")
  • F.column("column_name")

列を参照することは、select を実行することと同じではありません。列の「選択」とは、結果のデータセットに含める列を絞り込み、並べ替えることを指します。

目次

スキーマの取得

DataFrame.columns

すべての列名を Python のリストとして返します。

Copied!
1 columns = df.columns # ['age', 'name']

DataFrame.dtypes

すべての列名とそのデータ型をタプルのリストとして返します。

Copied!
1 dtypes = df.dtypes # [('age', 'int'), ('name', 'string')]

選択

DataFrame.select(*cols)

元の DataFrame の列の一部を含む、新しい DataFrame を返します。

たとえば、次の名前を持つ6列の DataFrame があるとします: id、first_name、last_name、phone_number、address、is_active_member

idfirst_namelast_namephone_numberzip_codeis_active_member
1JohnDoe(123) 456-789010014true
2JaneEyre(213) 555-123490007true
..................

DataFrame を加工して、必要な名前付き列(利用可能な列の一部)だけを含めたい場合があります。phone_number の1列だけを含むテーブルが必要だとします。

Copied!
1 df = df.select("phone_number")
phone_number
(123) 456-7890
(213) 555-1234
...

また、id、first_name、last_name だけが必要な場合もあります(同じ処理を実現する方法は少なくとも3つあります)。

  1. 列名を直接渡す方法:

    Copied!
    1 df = df.select("id", "first_name", "last_name")

    または、列のインスタンスを渡す方法:

    Copied!
    1 df = df.select(F.col("id"), F.col("first_name"), F.col("last_name"))
  2. 列名の配列を渡す方法:

    Copied!
    1 2 select_columns = ["id", "first_name", "last_name"] df = df.select(select_columns)
  3. 「アンパックした」配列を渡す方法:

    Copied!
    1 2 3 select_columns = ["id", "first_name", "last_name"] df = df.select(*select_columns) # 次と同じです: df = df.select("id", "first_name", "last_name")
    idfirst_namelast_name
    1JohnDoe
    2JaneEyre
    .........

    select_columns の前にある * は配列をアンパックし、機能上は #1 と同じ動作になるようにします(コメントを参照)。これにより、次のような処理ができます。

    Copied!
    1 2 3 select_columns = ["id", "first_name", "last_name"] return df.select(*select_columns, "phone_number") # 次と同じです: df = df.select("id", "first_name", "last_name", "phone_number")
    idfirst_namelast_namephone_number
    1JohnDoe(123) 456-7890
    2JaneEyre(213) 555-1234
    ............

出力データセットには、元の列の順序が保持されるのではなく、選択した列だけが、選択した順序で含まれることに注意してください。名前は一意で、大文字と小文字を区別し、選択元のデータセットに列としてすでに存在している必要があります。

このルールの例外として、新しい列を導出し、その列をすぐに選択することもできます。新しく導出した列には alias、つまり名前を付ける必要があります。

string1string2string3string4
firstsecondthirdFourth
onetwothreefour
Copied!
1 2 derived_column = F.concat_ws(":", F.col("string1"), F.col("string2")) return df.select("string3", derived_column.alias("derived"))
string3derived
thirdfirst
threeone

作成、更新

DataFrame.withColumn(name, column)

Copied!
1 new_df = old_df.withColumn("column_name", derived_column)
  • new_df:old_df のすべての列に加えて、new_column_name が追加された結果のデータフレームです。
  • old_df:新しい列を追加する対象のデータフレームです。
  • column_name:作成する列(old_df に存在しない場合)または更新する列(old_df にすでに存在する場合)の名前です。
  • derived_column:列を導出する式です。column_name(または列に付けた名前)のすべての行に適用されます。

既存の DataFrame に対して、withColumn メソッドを使用すると、新しい列を作成したり、既存の列を新しい値や変更した値で更新したりできます。これは、特に次の目的に役立ちます。

  1. 既存の値に基づく新しい値の導出

    Copied!
    1 df = df.withColumn("times_two", F.col("number") * 2) # times_two = number * 2
    Copied!
    1 df = df.withColumn("concat", F.concat(F.col("string1"), F.col("string2")))
  2. あるデータ型から別のデータ型への値のキャスト

    Copied!
    1 2 # `start_timestamp` を DateType にキャストし、新しい値を `start_date` に保存します df = df.withColumn("start_date", F.col("start_timestamp").cast("date"))
  3. 列の更新

    Copied!
    1 2 # 列 `string` を、その値をすべて小文字にしたものに更新します df = df.withColumn("string", F.lower(F.col("string")))

名前の変更、エイリアス

DataFrame.withColumnRenamed(name, rename)

列の名前を変更するには、.withColumnRenamed() を使用します。

Copied!
1 df = df.withColumnRenamed("old_name", "new_name")

列の名前を変更する処理は、次のように捉えることもできます。これは、PySpark が変換ステートメントをどのように最適化するかを理解する手がかりになります。

Copied!
1 df = df.withColumn("new_name", F.col("old_name")).drop("old_name")

ただし、withColumn を使わずに新しい列を導出し、その列に名前を付ける必要がある場合もあります。このような場合に役立つのが alias(またはそのメソッドの別名である name)です。以下に使用例を示します。

Copied!
1 2 3 4 df = df.select(derived_column.alias("new_name")) df = df.select(derived_column.name("new_name")) # .alias("new_name") と同じです df = df.groupBy("group") \ .agg(F.sum("number").alias("sum_of_numbers"), F.count("*").alias("count"))

複数の列の名前を一度に変更することもできます。

Copied!
1 2 3 4 5 6 7 renames = { "column": "column_renamed", "data": "data_renamed", } for colname, rename in renames.items(): df = df.withColumnRenamed(colname, rename)

削除

DataFrame.drop(*cols)

元の DataFrame から指定した列を削除し、残りの列を含む新しい DataFrame を返します。(指定した列名がスキーマに含まれていない場合、この処理に失敗します。)

列を削除するには、直接的な方法と間接的な方法の2つがあります。間接的な方法では、select を使用して、保持したい列の一部を選択します。直接的な方法では、drop を使用して、削除したい列の一部を指定します。どちらも使用時の構文は似ていますが、ここでは順序は重要ではありません。以下に例を示します。

idfirst_namelast_namephone_numberzip_codeis_active_member
1JohnDoe(123) 456-789010014true
2JaneEyre(213) 555-123490007true
..................

phone_number の1列だけを削除するとします。

Copied!
1 df = df.drop("phone_number")
idfirst_namelast_namezip_codeis_active_member
1JohnDoe10014true
2JaneEyre90007true
...............

また、id、first_name、last_name を削除したい場合もあります(同じ処理を実現する方法は少なくとも3つあります)。

  1. 列名を直接渡す方法:

    Copied!
    1 df = df.drop("id", "first_name", "last_name")

    または

    Copied!
    1 df = df.drop(F.col("id"), F.col("first_name"), F.col("last_name"))
  2. 配列を渡す方法:

    Copied!
    1 2 drop_columns = ["id", "first_name", "last_name"] df = df.drop(drop_columns)
  3. 「アンパックした」配列を渡す方法:

    Copied!
    1 2 3 drop_columns = ["id", "first_name", "last_name"] df = df.drop(*drop_columns) # 次と同じです: df = df.drop("id", "first_name", "last_name")
    phone_numberzip_codeis_active_member
    (123) 456-789010014true
    (213) 555-123490007true
    .........

    drop_columns の前にある * は配列をアンパックし、機能上は #1 と同じ動作になるようにします(コメントを参照)。これにより、次のような処理ができます。

    Copied!
    1 2 3 drop_columns = ["id", "first_name", "last_name"] df = df.drop(*drop_columns, "phone_number") # 次と同じです: df = df.drop("id", "first_name", "last_name", "phone_number")
    zip_codeis_active_member
    10014true
    90007true
    ......

キャスト

Column.cast(type)

存在するデータ型は次のとおりです:NullType、StringType、BinaryType、BooleanType、DateType、TimestampType、DecimalType、DoubleType、FloatType、ByteType、IntegerType、LongType、ShortType、ArrayType、MapType、StructType、StructField

一般に、列に対して cast メソッドを使うと、ほとんどのデータ型を別のデータ型に変換できます。

Copied!
1 2 3 from pyspark.sql.types import StringType df.select(df.age.cast(StringType()).alias("age")) # df.age は IntegerType() であると仮定します

または

Copied!
1 2 df.select(df.age.cast("string").alias("age")) # 実質的に StringType() を使う場合と同じです。
age
"2"
"5"

キャストすると、実質的に新しい派生列が作成され、その列に対して直接 select、withColumn、filter などを実行できます。「ダウンキャスト」と「アップキャスト」の概念は PySpark にも適用されるため、元のデータ型に格納されていたより詳細な情報が失われたり、不要な情報が加わったりする場合があります。

When、otherwise

F.when(condition, value).otherwise(value2)

パラメーター:

  • condition - ブール値の Column 式
  • value - リテラル値または Column 式

value または value2 パラメーターと同一の列の式として評価されます。Column.otherwise() が呼び出されない場合、一致しない条件に対しては None (null) の列の式が返されます。

when、otherwise 演算子は、新しい列の値を計算する if-else 文に相当します。基本的な使い方は次のとおりです。

Copied!
1 2 3 4 5 # CASE WHEN (age >= 21) THEN true ELSE false END at_least_21 = F.when(F.col("age") >= 21, True).otherwise(False) # CASE WHEN (last_name != "") THEN last_name ELSE null last_name = F.when(F.col("last_name") != "", F.col("last_name")).otherwise(None) df = df.select(at_least_21.alias("at_least_21"), last_name.alias("last_name"))

when 文は、必要な回数だけ連結することもできます。

Copied!
1 switch = F.when(F.col("age") >= 35, "A").when(F.col("age") >= 21, "B").otherwise("C")

これらの評価結果は、列に割り当てたり、フィルターで使用したりできます。

Copied!
1 2 df = df.withColumn("switch", switch) # switch=A, B, または C df = df.where(~F.isnull(last_name)) # 空の文字列を null 値として評価した後、last_name が null ではない行に絞り込みます