Skip to content

DNET_SparkSQL

nishi_74322014 edited this page Sep 11, 2026 · 1 revision

Spark SQL

概要

  • SQLで記述できる構造化データの処理の分散処理フレームワーク

  • Spark SQLとDataFrame APIは、DataFrameで同じことができる。

詳細

機能

  • バックエンドで Apache Sparkが動作する、

  • Spark Core(Apache Sparkの該当節を参照)より上位のコンポーネント

  • Apache Hiveとの関連

    • 先行する Apache Hive と互換性を保ちながら発展中
    • Apache Hive で培われた SQLによる分散処理を利用

ドメイン固有言語 (DSL) > SQL言語

  • RDD(Apache Sparkの該当節を参照)の処理を
    クエリ・エンジン経由で実行プラン生成して処理。

    • 基本的にはDataFrameというオンメモリのテーブルに対し、
      LINQ的に処理を行うことが出来る。
    • 強く型付けされたデータセットも完全にサポート(DataFrameは型無しAPI)
  • サポート

  • クエリ オプティマイザ

    • Catalyst Optimizer

    • サポート

      • ルールベース
      • コストベース
    • コストベースの 4 フェーズ

      • 分析
      • ロジックの最適化
      • 物理的プランニング
      • コード生成
  • 遅延評価(.NET for Apache Sparkチュートリアルの該当節を参照)

    • 実行計画は有向グラフ(DAG)として表現される。
    • エグゼキュータでDAGが生成され、
    • DAGを用いた遅延評価で、一回、処理が解析される。
    • RDDを消費するAPIが呼ばれるとジョブが作られる。
    • ドライバはエグゼキュータにジョブを転送する。

DataFrameと言う構造化データ

  • テーブルのようなデータ構造をもった分散処理用データセット

    • 同名のRのデータ構造をモデルとした
      表形式データ(型と名前を持った列のテーブル)

    • 作業性のサポート向上を目的としている。

  • 今後は、RDD(Apache Sparkの該当節を参照)ではなく
    DataFrameが主流になって行く。

    • RDDを上回る抽象化が提供される。
    • Catalyst オプティマイザでRDDの5~20倍、性能向上。

構文

LINQ

「クエリ構文」的記法

  • SparkSessionのsqlメソッドで入力元にクエリを発行し、結果をDataFrameで返す。
  • 若しくは、DataFrameにテーブル名を指定してSQL中で使用することでDataFrameを操作する。

「メソッド構文」的記法

コチラは、DataFrame操作のみの場合。

  • 正確にはDataFrame APIと呼ぶらしい。
  • DataFrameをFluent API的なメソッド構文で操作する。

操作

基本(操作)

PySparkの例(PySparkの該当節を参照)

ユーザー定義関数(UDF)

集計

基本(集計)

PySparkの例(PySparkの該当節を参照)

ユーザー定義集計関数(UDAF)

Apache Spark 2.3 以降では Pandas UDF を使うことで
UDAF に相当する処理を書くことができるようになっている。

構造化ストリーミング

Spark Structured Streaming

Python

PySparkの該当節を参照。

.NET

構造化ストリーミングを参照。

チュートリアル

Scala

Apache Sparkチュートリアルの該当節を参照。

PySpark

PySparkの該当節を参照。

ファースト・ステップ的なモノ

PySparkの該当節を参照。

DataFrameに対する様々な操作

PySparkの該当節を参照。

構造化ストリーミングでウィンドウ操作

PySparkの該当節を参照。

様々な入力元・出力先の例

PySparkの該当節を参照。

.NET for Apache Spark

Get started in 10 minutes

.NET for Apache Sparkチュートリアルの該当節を参照。

  • データ
Hello World
This .NET app uses .NET for Apache Spark
This .NET app counts words with Apache Spark
  • コード
// Create Spark session
SparkSession spark =
  SparkSession
    .Builder()
    .AppName("word_count_sample")
    .GetOrCreate();

// Create initial DataFrame
string filePath = args[0];
DataFrame dataFrame = spark.Read().Text(filePath);

// Count words
DataFrame words =
  dataFrame
    .Select(Split(Col("value")," ").Alias("words"))
    .Select(Explode(Col("words")).Alias("word"))
    .GroupBy("word")
    .Count()
    .OrderBy(Col("count").Desc());

// Display results
words.Show();

// Stop Spark session
spark.Stop();
  • 説明

    • value列(既定の列名)を、半角スペースで分割して配列化、wordsと言うaliasの列名を設定する。

      Select(Split(Col("value"), " ")).Alias("words"))
    • Explodeで、前述の結果セットのwords列セル中の配列を行に展開し、

      .Select(Explode(Col("words"))
    • 前述の結果セットをword列で行数を集計し、
      集計列(列名は"count"になる)をDescでOrderByする。

      .GroupBy("word").Count().OrderBy(Col("count").Desc())
    • 参考
      pandas documentation

バッチ処理

.NET for Apache Sparkチュートリアルの該当節を参照。

id INT url STRING owner_id INT name STRING descriptor STRING language STRING created_at STRING forked_from INT deleted STRING updated_at STRING
  • コード
// Create Spark session
SparkSession spark = SparkSession
  .Builder()
  .AppName("GitHub and Spark Batch")
  .GetOrCreate();

// Create initial DataFrame
DataFrame projectsDf = spark
  .Read()
  .Schema("id INT, url STRING, owner_id INT, " +
  "name STRING, descriptor STRING, language STRING, " +
  "created_at STRING, forked_from INT, deleted STRING, " +
  "updated_at STRING")
  .Csv(filePath);

// Display results
projectsDf.Show();

// Drop any rows with NA values
DataFrameNaFunctions dropEmptyProjects = projectsDf.Na();
DataFrame cleanedProjects = dropEmptyProjects.Drop("any");

// Remove unnecessary columns
cleanedProjects = cleanedProjects.Drop("id", "url", "owner_id");

// Display results
cleanedProjects.Show();

// Average number of times each language has been forked
DataFrame groupedDF = cleanedProjects
  .GroupBy("language")
  .Agg(Avg(cleanedProjects["forked_from"]));

// Sort by most forked languages first & Display results
groupedDF.OrderBy(Desc("avg(forked_from)")).Show();

// Defines a UDF that determines if a date is greater than a specified date
spark.Udf().Register<string, bool>(
  "MyUDF",
  (date) => DateTime.TryParse(date, out DateTime convertedDate) &&
    (convertedDate > referenceDate));

// Use UDF to add columns to the generated TempView
cleanedProjects.CreateOrReplaceTempView("dateView");
DataFrame dateDf = spark.Sql(
  "SELECT *, MyUDF(dateView.updated_at) AS datebefore FROM dateView");

// Display results
dateDf.Show();

// Stop Spark session
spark.Stop();
  • 説明

    • 初めに、

      • 任意のNULLまたはNaN値を含む行を削除
      • 不要な列の削除("id", "url", "owner_id")
    • 次に、各言語がフォークされた回数の平均値を求める。

      • 言語毎に集計

        .GroupBy("language")
      • 集計方法はforked_from列の平均値

        .Agg(Avg(cleanedProjects["forked_from"]));
      • 最後に集計列(列名は"avg(forked_from)"になる)をDescでOrderByする。

    • 最後に、datebefore列(updated_atが指定の日付より大きいか?)を追加。

構造化ストリーミング(.NET)

.NET for Apache Sparkチュートリアルの該当節を参照。

  • データ
    ncからレスポンスされる文字列。

  • コード

// Create Spark session
SparkSession spark = SparkSession
  .Builder()
  .AppName("Streaming example with a UDF")
  .GetOrCreate();

// Create initial DataFrame
DataFrame lines = spark
  .ReadStream()
  .Format("socket")
  .Option("host", hostname)
  .Option("port", port)
  .Load();

// UDF to produce an array
// Array includes:
// 1) original string
// 2) original string + length of original string
Func<Column, Column> udfArray =
  Udf<string, string[]>((str) => new string[] { str, $"{str} {str.Length}" });

// Explode array to rows
DataFrame arrayDF = lines.Select(Explode(udfArray(lines["value"])));

// Process and display each incoming line
StreamingQuery query = arrayDF
  .WriteStream()
  .Format("console")
  .Start();

query.AwaitTermination();

※ WriteStreamでStartして取得したQueryにAwaitTerminationして終了。

  • 説明

    • 来た文字列を、{文字列, 文字列長}の
      文字列配列にするユーザー定義関数(UDF)を登録
    • メソッド構文を使用してUDFを実行
    • UDFの返す配列をExplodeで行に展開して表示。
    • 構造化ストリーミングはループを書かずして、
      ストリーミング処理をループ的に処理することができる。

ML.NETでの感情分析

.NET for Apache Sparkチュートリアルの該当節を参照。

参考

DevelopersIO

Hadoop Advent Calendar 2016

Learn | Microsoft Docs

マイクロソフト系技術情報 Wiki

移行メモ

  • 「Get started in 10 minutes」の説明「wordと言うaliasの列名を設定する」は、 直後のコード(.Alias("words"))に合わせ「wordsと言うaliasの列名」に正した。
  • マイクロソフト系技術情報 Wiki(techinfoofmicrosofttech.osscons.jp)への リンクは、移行済みの MS_DotNetForApacheSparkTutorial / MS_LINQ / MS_AzureDatabricksTutorial に張り替えた。

Tags: 移行, Apache Spark, Spark SQL, DataFrame, 分散処理, LINQ, .NET for Apache Spark

NetDevInfraWiki

マイクロソフト系技術情報 Wiki
Open 棟梁 Wiki

(未着手)

開発基盤部会 Wiki

移行管理: DONETODO

Clone this wiki locally