-
Notifications
You must be signed in to change notification settings - Fork 0
MS_DotNetForApacheSparkDataAccess
- 戻る(.NET for Apache Spark)
- .NET for Apache Sparkチュートリアル
- .NET for Apache Sparkのデータ接続
- .NET for Apache SparkのSQL
.NET で Spark の実装ができる for Apache Spark。
Spark の仕様なので、Azure HDInsight・Azure Databricksの
差異はない。
ローカル・ファイルを読むこともできる。
マウントされたストレージに直接アップロードしてソレを読む。
Databricks CLI を使用してアップロードしてソレを読む。
下記に対する 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 を使う版で、こちらを使う。
- SparkSession の
ReadStreamメソッドで DataFrame を読み込む。 - Loop ではなく StreamingQuery を使用して実装する。
.Format("socket")
- ホスト名
.Option("host", hostname) - ポート番号
.Option("port", port)
Azure Databricksチュートリアルが参考になる。
正確には、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()が返すのはDataStreamWriter、Start()が返すのは
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の指定は通常不要(指定すると警告や競合の原因になる)。
サーバーにはならない
- 使い方ガイド > データーへの接続
- Azure Data Lake Storage Gen 2 または WASB アカウントに接続する
https://docs.microsoft.com/ja-jp/dotnet/spark/how-to-guides/connect-to-azure-storage - .NET for Apache Spark を Azure Event Hubs に接続する
https://docs.microsoft.com/ja-jp/dotnet/spark/how-to-guides/connect-to-event-hub - .NET for Apache Spark を MongoDB に接続する
https://docs.microsoft.com/ja-jp/dotnet/spark/how-to-guides/connect-to-mongo-db - .NET for Apache Spark を SQL Server に接続する
https://docs.microsoft.com/ja-jp/dotnet/spark/how-to-guides/connect-to-sql-server
- Azure Data Lake Storage Gen 2 または WASB アカウントに接続する
- .NET for Apache Spark
- .NET for Apache Sparkチュートリアル
- Azure HDInsight / Azure Databricks
- Azure Event Hubs
Tags: クラウド, Azure, .NET開発, .NET Core, .NET Standard
このWikiは「Open棟梁Project」,「OSSコンソーシアム 開発基盤部会」によって運営されています。