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!
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();
}
}
}
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
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
@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-submitcommand. The--packagesoption 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!
Most helpful comment
@Ning-Guo the following can be found in your stack trace:
In the Structured Streaming + Kafka Integration Guide there is a Deploying section that contains an example on how to construct your
spark-submitcommand. The--packagesoption will need to be added. You will need to use the kafka package that matches your spark version.