Skip to content

MS_DotNetForApacheSparkDataAccess

nishi_74322014 edited this page Sep 1, 2026 · 1 revision

.NET for Apache Sparkのデータ接続

概要

.NET で Spark の実装ができる for Apache Spark。

詳細

Spark の仕様なので、Azure HDInsightAzure Databricks
差異はない。

ファイルの場合

ローカル

ローカル・ファイルを読むこともできる。

Azure HDInsight の場合

マウントされたストレージに直接アップロードしてソレを読む。

Azure Databricks の場合

Databricks CLI を使用してアップロードしてソレを読む。

WASB or ABFS

下記に対する URL を使用することも出来る。

  • WASB[S]: Azure storage.
    wasb://<container_name>@<storage_account_name>.blob.core.windows.net

  • ABFS[S]: Azure Data Lake Storage Gen2.
    abfs://<file_system>@<account_name>.dfs.core.windows.net

補足(WASB は非推奨): WASB(Blob 用のドライバ)は既に非推奨で、
現在は **ABFS(Azure Data Lake Storage Gen2)**を使うのが標準である。
また、末尾に S が付く wasbs:// / abfss://
TLS を使う版で、こちらを使う。

Clientの場合

TCP/IPソケット

.Format("socket")

  • ホスト名
    .Option("host", hostname)
  • ポート番号
    .Option("port", port)

Azure Event Hubs

Azure Databricksチュートリアルが参考になる。

Apache Kafka

正確には、Azure Event Hubs の Kafka エンドポイント
Azure Event Hubsチュートリアル

  • 受信処理(DataFrame の ReadStream)の結果を DataFrame.Show() で出力する。

  • なお、送信処理(DataFrame の WriteStream)もサポートしている。

  • 送受信処理の実装

    • パラメタの初期化(共通)
      SAS トークンを使用する場合の設定例

      string BOOTSTRAP_SERVERS = "hostname:9093"; // 9093 is the port used to communicate with Event Hubs, see [troubleshooting guide](https://docs.microsoft.com/azure/event-hubs/troubleshooting-guide)
      string EH_SASL = "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"$ConnectionString\" password=\"<CONNECTION_STRING>\";"; // Connection string obtained from Step 1
    • ReadStream

      SparkSession spark = SparkSession
          .Builder()
          .AppName("Connect Event Hub")
          .GetOrCreate();
      
      DataFrame df = spark
          .ReadStream()
          .Format("kafka")
          .Option("kafka.bootstrap.servers", BOOTSTRAP_SERVERS)
          .Option("subscribe", "spark-test")
          .Option("kafka.sasl.mechanism", "PLAIN")
          .Option("kafka.security.protocol", "SASL_SSL")
          .Option("kafka.sasl.jaas.config", EH_SASL)
          .Option("kafka.request.timeout.ms", "60000")
          .Option("kafka.session.timeout.ms", "60000")
          .Option("failOnDataLoss", "false")
          .Load();
      
      DataFrame dfWrite = df
          .WriteStream()
          .OutputMode("append")
          .Format("console")
          .Start();
    • WriteStream

      df.WriteStream()
          .Format("kafka")
          .Option("topic", topics)
          .Option("kafka.bootstrap.servers", BOOTSTRAP_SERVERS)
          .Option("kafka.sasl.mechanism", "PLAIN")
          .Option("kafka.security.protocol", "SASL_SSL")
          .Option("kafka.sasl.jaas.config", EH_SASL)
          .Option("checkpointLocation", "./checkpoint")
          .Start();

移行メモ(型の不一致): ReadStream の例で
DataFrame dfWrite = df.WriteStream()...Start(); となっているが、
WriteStream() が返すのは DataStreamWriterStart() が返すのは
StreamingQuery であり、DataFrame ではない。
元コードのまま残すが、実装時は StreamingQuery で受けること
(直後の「チュートリアルの書き換え」でも StreamingQuery を使う旨が
リンク先で述べられている)。

  • チュートリアルの書き換え
    .NET for Apache Spark の構造化ストリーミングのチュートリアル
    受信部を書き換えて実行してみる。

    • 書換

      // Create initial DataFrame
      string BOOTSTRAP_SERVERS = "xxxx.servicebus.windows.net:9093";
      string CONNECTION_STRING = "Endpoint=sb://xxxx.servicebus.windows.net/;SharedAccessKeyName=RootManageSharedAccessKey;SharedAccessKey=xxxx";
      string EH_SASL = $"org.apache.kafka.common.security.plain.PlainLoginModule required username=\"$ConnectionString\" password=\"{CONNECTION_STRING}\";";
      DataFrame lines = spark
          .ReadStream()
          .Format("kafka")
          .Option("subscribe", "test")
          .Option("kafka.bootstrap.servers", BOOTSTRAP_SERVERS)
          .Option("kafka.sasl.mechanism", "PLAIN")
          .Option("kafka.security.protocol", "SASL_SSL")
          .Option("kafka.sasl.jaas.config", EH_SASL)
          .Option("kafka.request.timeout.ms", "60000")
          .Option("kafka.session.timeout.ms", "60000")
          .Option("kafka.group.id", "$Default")
          .Option("failOnDataLoss", "false")
          .Load();
    • 実行

      >spark-submit ^
      --class org.apache.spark.deploy.dotnet.DotnetRunner ^
      --master local ^
      microsoft-spark-3-0_2.12-2.0.0.jar ^
      dotnet mySparkStreamingApp.dll
      

      ※ エラー発生中(調査中)

補足(Kafka 接続で詰まりやすい点): 上記が「調査中」で止まっている件について、
Event Hubs の Kafka エンドポイントに繋ぐ場合によく問題になるのは次の 3 点である。

  • Kafka コネクタの JAR が無い
    spark-sql-kafka-0-10--packages で指定する必要がある
    .Format("kafka") は Spark 本体には含まれない)。
  • Event Hubs の SKU
    Kafka エンドポイントは Basic では使えない(Standard 以上)。
  • kafka.group.id
    Spark の Structured Streaming は自前でオフセットを管理するため、
    group.id の指定は通常不要(指定すると警告や競合の原因になる)。

Serverの場合

サーバーにはならない

参考

microsoft.com

Microsoft Docs > .NET for Apache Spark ガイド

.NET for Apache Spark

本 Wiki 内


Tags: クラウド, Azure, .NET開発, .NET Core, .NET Standard

NetDevInfraWiki

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

(未着手)

開発基盤部会 Wiki

移行管理: DONETODO

Clone this wiki locally