Spark: [ERROR]: got error on example "csharp streaming StructuredKafkaWordCount.cs"

Created on 18 Nov 2019  路  2Comments  路  Source: dotnet/spark

Describe the bug
I would like to use the dotnet spark to consume my kafka streaming queue here, and when I run my application I got below errors, could someone kindly help me out of this problem? Thanks a lot!

  1. Csharp code :

using System;
using Microsoft.Spark.Sql;
using static Microsoft.Spark.Sql.Functions;
namespace HelloSpark
{
class Program
{
static void Main(string[] args)
{
SparkSession spark = SparkSession.Builder().AppName("my spark kafka intergration").GetOrCreate();
string server = "xxxxx"; //my kafka server name
string type = "subscribe";
string topics = "xxxx"; //my kafka topic name
DataFrame lines = spark.ReadStream().Format("kafka").Option("kafka.bootstrap.servers", server).Option(type, topics).Load().SelectExpr("CAST(value AS STRING)");
DataFrame words = lines.Select(Explode(Split(lines["value"], " ")).Alias("word"));
DataFrame wordcounts = words.GroupBy("word").Count();
Microsoft.Spark.Sql.Streaming.StreamingQuery query = wordcounts.WriteStream().OutputMode("complete").Format("console").Start();
query.AwaitTermination();
}
}
}

  1. in command promote
    %SPARK_HOME%\bin\spark-submit --class org.apache.spark.deploy.dotnet.DotnetRunner --master local bin\Debug\netcoreapp3.0\microsoft-spark-2.4.x-0.6.0.jar dotnet bin\Debug\netcoreapp3.0\HelloSpark.dll

  2. See error
    19/11/19 01:40:28 ERROR DotnetBackendHandler: methods:
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.slf4j.Logger org.apache.spark.sql.streaming.DataStreamReader.log()
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.apache.spark.sql.streaming.DataStreamReader org.apache.spark.sql.streaming.DataStreamReader.format(java.lang.String)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.apache.spark.sql.Dataset org.apache.spark.sql.streaming.DataStreamReader.load()
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.apache.spark.sql.Dataset org.apache.spark.sql.streaming.DataStreamReader.load(java.lang.String)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public java.lang.String org.apache.spark.sql.streaming.DataStreamReader.logName()
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public void org.apache.spark.sql.streaming.DataStreamReader.logInfo(scala.Function0,java.lang.Throwable)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public void org.apache.spark.sql.streaming.DataStreamReader.logInfo(scala.Function0)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public void org.apache.spark.sql.streaming.DataStreamReader.logTrace(scala.Function0)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public void org.apache.spark.sql.streaming.DataStreamReader.logTrace(scala.Function0,java.lang.Throwable)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public void org.apache.spark.sql.streaming.DataStreamReader.logWarning(scala.Function0,java.lang.Throwable)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public void org.apache.spark.sql.streaming.DataStreamReader.logWarning(scala.Function0)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public boolean org.apache.spark.sql.streaming.DataStreamReader.isTraceEnabled()
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public void org.apache.spark.sql.streaming.DataStreamReader.logError(scala.Function0)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public void org.apache.spark.sql.streaming.DataStreamReader.logError(scala.Function0,java.lang.Throwable)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public void org.apache.spark.sql.streaming.DataStreamReader.logDebug(scala.Function0,java.lang.Throwable)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public void org.apache.spark.sql.streaming.DataStreamReader.logDebug(scala.Function0)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public boolean org.apache.spark.sql.streaming.DataStreamReader.initializeLogIfNecessary$default$2()
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public void org.apache.spark.sql.streaming.DataStreamReader.org$apache$spark$internal$Logging$$log__$eq(org.slf4j.Logger)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.slf4j.Logger org.apache.spark.sql.streaming.DataStreamReader.org$apache$spark$internal$Logging$$log_()
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public void org.apache.spark.sql.streaming.DataStreamReader.initializeLogIfNecessary(boolean)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public boolean org.apache.spark.sql.streaming.DataStreamReader.initializeLogIfNecessary(boolean,boolean)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.apache.spark.sql.streaming.DataStreamReader org.apache.spark.sql.streaming.DataStreamReader.options(scala.collection.Map)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.apache.spark.sql.streaming.DataStreamReader org.apache.spark.sql.streaming.DataStreamReader.options(java.util.Map)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.apache.spark.sql.streaming.DataStreamReader org.apache.spark.sql.streaming.DataStreamReader.option(java.lang.String,java.lang.String)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.apache.spark.sql.streaming.DataStreamReader org.apache.spark.sql.streaming.DataStreamReader.option(java.lang.String,boolean)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.apache.spark.sql.streaming.DataStreamReader org.apache.spark.sql.streaming.DataStreamReader.option(java.lang.String,long)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.apache.spark.sql.streaming.DataStreamReader org.apache.spark.sql.streaming.DataStreamReader.option(java.lang.String,double)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.apache.spark.sql.Dataset org.apache.spark.sql.streaming.DataStreamReader.text(java.lang.String)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.apache.spark.sql.streaming.DataStreamReader org.apache.spark.sql.streaming.DataStreamReader.schema(org.apache.spark.sql.types.StructType)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.apache.spark.sql.streaming.DataStreamReader org.apache.spark.sql.streaming.DataStreamReader.schema(java.lang.String)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.apache.spark.sql.Dataset org.apache.spark.sql.streaming.DataStreamReader.orc(java.lang.String)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.apache.spark.sql.Dataset org.apache.spark.sql.streaming.DataStreamReader.csv(java.lang.String)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.apache.spark.sql.Dataset org.apache.spark.sql.streaming.DataStreamReader.textFile(java.lang.String)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.apache.spark.sql.Dataset org.apache.spark.sql.streaming.DataStreamReader.parquet(java.lang.String)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public org.apache.spark.sql.Dataset org.apache.spark.sql.streaming.DataStreamReader.json(java.lang.String)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public final void java.lang.Object.wait() throws java.lang.InterruptedException
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public final void java.lang.Object.wait(long,int) throws java.lang.InterruptedException
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public final native void java.lang.Object.wait(long) throws java.lang.InterruptedException
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public boolean java.lang.Object.equals(java.lang.Object)
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public java.lang.String java.lang.Object.toString()
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public native int java.lang.Object.hashCode()
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public final native java.lang.Class java.lang.Object.getClass()
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public final native void java.lang.Object.notify()
    19/11/19 01:40:28 ERROR DotnetBackendHandler: public final native void java.lang.Object.notifyAll()
    19/11/19 01:40:28 ERROR DotnetBackendHandler: args:
    [2019-11-18T17:40:28.9107466Z] [MININT-RRGK6M1] [Error] [JvmBridge] JVM method execution failed: Nonstatic method load failed for class 6 when called with no arguments
    [2019-11-18T17:40:28.9108882Z] [MININT-RRGK6M1] [Error] [JvmBridge] org.apache.spark.sql.AnalysisException: Failed to find data source: kafka. Please deploy the application as per the deployment section of "Structured Streaming + Kafka Integration Guide".;
    at org.apache.spark.sql.execution.datasources.DataSource$.lookupDataSource(DataSource.scala:652)
    at org.apache.spark.sql.streaming.DataStreamReader.load(DataStreamReader.scala:161)
    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at sun.reflect.NativeMethodAccessorImpl.invoke(Unknown Source)
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(Unknown Source)
    at java.lang.reflect.Method.invoke(Unknown Source)
    at org.apache.spark.api.dotnet.DotnetBackendHandler.handleMethodCall(DotnetBackendHandler.scala:162)
    at org.apache.spark.api.dotnet.DotnetBackendHandler.handleBackendRequest(DotnetBackendHandler.scala:102)
    at org.apache.spark.api.dotnet.DotnetBackendHandler.channelRead0(DotnetBackendHandler.scala:29)
    at org.apache.spark.api.dotnet.DotnetBackendHandler.channelRead0(DotnetBackendHandler.scala:24)
    at io.netty.channel.SimpleChannelInboundHandler.channelRead(SimpleChannelInboundHandler.java:105)
    at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:362)
    at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:348)
    at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:340)
    at io.netty.handler.codec.MessageToMessageDecoder.channelRead(MessageToMessageDecoder.java:102)
    at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:362)
    at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:348)
    at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:340)
    at io.netty.handler.codec.ByteToMessageDecoder.fireChannelRead(ByteToMessageDecoder.java:310)
    at io.netty.handler.codec.ByteToMessageDecoder.channelRead(ByteToMessageDecoder.java:284)
    at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:362)
    at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:348)
    at io.netty.channel.AbstractChannelHandlerContext.fireChannelRead(AbstractChannelHandlerContext.java:340)
    at io.netty.channel.DefaultChannelPipeline$HeadContext.channelRead(DefaultChannelPipeline.java:1359)
    at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:362)
    at io.netty.channel.AbstractChannelHandlerContext.invokeChannelRead(AbstractChannelHandlerContext.java:348)
    at io.netty.channel.DefaultChannelPipeline.fireChannelRead(DefaultChannelPipeline.java:935)
    at io.netty.channel.nio.AbstractNioByteChannel$NioByteUnsafe.read(AbstractNioByteChannel.java:138)
    at io.netty.channel.nio.NioEventLoop.processSelectedKey(NioEventLoop.java:645)
    at io.netty.channel.nio.NioEventLoop.processSelectedKeysOptimized(NioEventLoop.java:580)
    at io.netty.channel.nio.NioEventLoop.processSelectedKeys(NioEventLoop.java:497)
    at io.netty.channel.nio.NioEventLoop.run(NioEventLoop.java:459)
    at io.netty.util.concurrent.SingleThreadEventExecutor$5.run(SingleThreadEventExecutor.java:858)
    at io.netty.util.concurrent.DefaultThreadFactory$DefaultRunnableDecorator.run(DefaultThreadFactory.java:138)
    at java.lang.Thread.run(Unknown Source)

on windows 10

question

Most helpful comment

@Ning-Guo the following can be found in your stack trace:

[2019-11-18T17:40:28.9108882Z] [MININT-RRGK6M1] [Error] [JvmBridge] org.apache.spark.sql.AnalysisException: Failed to find data source: kafka. Please deploy the application as per the deployment section of "Structured Streaming + Kafka Integration Guide".;

In the Structured Streaming + Kafka Integration Guide there is a Deploying section that contains an example on how to construct your spark-submit command. The --packages option will need to be added. You will need to use the kafka package that matches your spark version.

All 2 comments

@Ning-Guo the following can be found in your stack trace:

[2019-11-18T17:40:28.9108882Z] [MININT-RRGK6M1] [Error] [JvmBridge] org.apache.spark.sql.AnalysisException: Failed to find data source: kafka. Please deploy the application as per the deployment section of "Structured Streaming + Kafka Integration Guide".;

In the Structured Streaming + Kafka Integration Guide there is a Deploying section that contains an example on how to construct your spark-submit command. The --packages option will need to be added. You will need to use the kafka package that matches your spark version.

@Ning-Guo the following can be found in your stack trace:

[2019-11-18T17:40:28.9108882Z] [MININT-RRGK6M1] [Error] [JvmBridge] org.apache.spark.sql.AnalysisException: Failed to find data source: kafka. Please deploy the application as per the deployment section of "Structured Streaming + Kafka Integration Guide".;

In the Structured Streaming + Kafka Integration Guide there is a Deploying section that contains an example on how to construct your spark-submit command. The --packages option will need to be added. You will need to use the kafka package that matches your spark version.

Thank a lot for your reply. Yes, this error was caused by the kafka streaming source. There is nothing wrong with the sample code here. Thanks a lot!

Was this page helpful?
0 / 5 - 0 ratings