Data Science WorkspaceのSparkを使用したデータへのアクセス

NOTE
Data Science Workspaceは購入できなくなりました。
このドキュメントは、Data Science Workspaceの使用権限を以前に持つ既存のお客様向けです。

次のドキュメントでは、Data Science Workspaceで使用するSparkを使用してデータにアクセスする方法の例を示します。 JupyterLab ノートブックを使用したデータへのアクセスについて詳しくは、JupyterLab ノートブックデータアクセス ​のドキュメントを参照してください。

はじめに

Sparkを使用するには、SparkSessionに追加する必要があるパフォーマンスの最適化が必要です。 さらに、後でデータセットに読み取りおよび書き込みを行うためにconfigPropertiesを設定することもできます。

import com.adobe.platform.ml.config.ConfigProperties
import com.adobe.platform.query.QSOption
import org.apache.spark.sql.{DataFrame, SparkSession}

Class Helper {

 /**
   *
   * @param configProperties - Configuration Properties map
   * @param sparkSession     - SparkSession
   * @return                 - DataFrame which is loaded for training
   */

   def load_dataset(configProperties: ConfigProperties, sparkSession: SparkSession, taskId: String): DataFrame = {
            // Read the configs
            val userToken: String = sparkSession.sparkContext.getConf.get("ML_FRAMEWORK_IMS_TOKEN", "").toString
            val orgId: String = sparkSession.sparkContext.getConf.get("ML_FRAMEWORK_IMS_ORG_ID", "").toString
            val apiKey: String = sparkSession.sparkContext.getConf.get("ML_FRAMEWORK_IMS_CLIENT_ID", "").toString
            val sandboxName: String = sparkSession.sparkContext.getConf.get("sandboxName", "").toString

   }
}

データセットの読み取り

Sparkを使用している間は、インタラクティブとバッチの2つの読み方にアクセスできます。

インタラクティブモードは、Query ServiceへのJava データベース接続(JDBC)接続を作成し、通常のJDBC ResultSetを通じて結果を取得します。この結果は、DataFrameに自動的に変換されます。 このモードは、組み込みのSpark メソッド spark.read.jdbc()と同様に機能します。 このモードは、小さなデータセットのみを対象としています。 データセットが500万行を超える場合は、バッチモードに切り替えることをお勧めします。

バッチモードでは、Query ServiceのCOPY コマンドを使用して、共有場所にParquet結果セットを生成します。 次に、これらのParquet ファイルをさらに処理できます。

インタラクティブモードでデータセットを読み取る例を次に示します。

  // Read the configs
    val userToken: String = sparkSession.sparkContext.getConf.get("ML_FRAMEWORK_IMS_TOKEN", "").toString
    val orgId: String = sparkSession.sparkContext.getConf.get("ML_FRAMEWORK_IMS_ORG_ID", "").toString
    val apiKey: String = sparkSession.sparkContext.getConf.get("ML_FRAMEWORK_IMS_CLIENT_ID", "").toString
    val sandboxName: String = sparkSession.sparkContext.getConf.get("sandboxName", "").toString

 val dataSetId: String = configProperties.get(taskId).getOrElse("")

    // Load the dataset
    var df = sparkSession.read.format(PLATFORM_SDK_PQS_PACKAGE)
      .option(QSOption.userToken, userToken)
      .option(QSOption.imsOrg, orgId)
      .option(QSOption.apiKey, apiKey)
      .option(QSOption.mode, "interactive")
      .option(QSOption.datasetId, dataSetId)
      .option(QSOption.sandboxName, sandboxName)
      .load()
    df.show()
    df
  }

同様に、バッチモードでデータセットを読み取る例を次に示します。

val df = sparkSession.read.format(PLATFORM_SDK_PQS_PACKAGE)
      .option(QSOption.userToken, userToken)
      .option(QSOption.imsOrg, orgId)
      .option(QSOption.apiKey, apiKey)
      .option(QSOption.mode, "batch")
      .option(QSOption.datasetId, dataSetId)
      .option(QSOption.sandboxName, sandboxName)
      .load()
    df.show()
    df

データセットから列を選択

df = df.select("column-a", "column-b").show()

DISTINCT句

DISTINCT句を使用すると、行/列レベルですべてのdistinct値を取得し、応答からすべての重複する値を削除できます。

distinct()関数の使用例を次に示します。

df = df.select("column-a", "column-b").distinct().show()

WHERE句

Spark SDKでは、SQL式を使用するか、条件を通じてフィルタリングするという2つのフィルタリング方法を使用できます。

これらのフィルタリング関数の使用例を次に示します。

SQL式

df.where("age > 15")

フィルター条件

df.where("age" > 15 || "name" = "Steve")

ORDER BY句

ORDER BY句を使用すると、受信した結果を特定の順序(昇順または降順)で指定した列で並べ替えることができます。 Spark SDKでは、これはsort()関数を使用して行われます。

sort()関数の使用例を次に示します。

df = df.sort($"column1", $"column2".desc)

LIMIT句

LIMIT句を使用すると、データセットから受信するレコードの数を制限できます。

limit()関数の使用例を次に示します。

df = df.limit(100)

データセットへの書き込み

configProperties マッピングを使用すると、QSOptionを使用してExperience Platformのデータセットに書き込むことができます。

val userToken: String = sparkSession.sparkContext.getConf.get("ML_FRAMEWORK_IMS_TOKEN", "").toString
val orgId: String = sparkSession.sparkContext.getConf.get("ML_FRAMEWORK_IMS_ORG_ID", "").toString
val apiKey: String = sparkSession.sparkContext.getConf.get("ML_FRAMEWORK_IMS_CLIENT_ID", "").toString
val sandboxName: String = sparkSession.sparkContext.getConf.get("sandboxName", "").toString

    df.write.format(PLATFORM_SDK_PQS_PACKAGE)
      .option(QSOption.userToken, userToken)
      .option(QSOption.imsOrg, orgId)
      .option(QSOption.apiKey, apiKey)
      .option(QSOption.datasetId, scoringResultsDataSetId)
      .option(QSOption.sandboxName, sandboxName)
      .save()

次の手順

Adobe Experience Platform Data Science Workspaceでは、上記のコードサンプルを使用してデータを読み書きするScala (Spark)のレシピサンプルを提供しています。 データへのアクセスにSparkを使用する方法について詳しくは、Data Science Workspace Scala GitHub Repositoryを参照してください。

recommendation-more-help
experience-platform-help-data-science-workspace