A PySpark job that calls spark.readStream.format("kafka") fails with Failed to find data source: kafka. Spark does not ship the Kafka source inside its core package. The Kafka connector, spark-sql-kafka-0-10, is a separate artifact, and the job fails when it is not on the classpath. The fix is to add the connector and its dependencies when the application is launched, not to change the code.
Why it happens
The Spark Kafka integration guide says that for Python applications you need to add the library and its dependencies when deploying the application. It says the same for experimenting in the Spark shell. Scala and Java applications link the artifact in their build instead, with the group id org.apache.spark and the artifact id spark-sql-kafka-0-10_2.13 for the current Scala build.
The fix
The guide’s deployment section adds the connector at launch with --packages. For a submitted application the command takes this shape, with the version set to your Spark version:
./bin/spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.13:4.2.0 ...
For the shell, the same --packages option works with spark-shell. The guide points to the Application Submission Guide for submitting applications with external dependencies.
Two details are worth checking before you launch.
- Match the connector version to your Spark version. The guide’s example uses the version of the Spark release it documents, so copy the version from the guide for your release, not from an older tutorial.
- Match the Scala build. The artifact name ends in a Scala version such as
_2.13, so take the full artifact name from the guide for your Spark release.
Check that the package loaded
Start the job with the package and run the read again. If the error is gone and you see KafkaSource in the query progress, the connector is on the classpath. The programming guide shows the source description in each progress report, for example KafkaSource[Subscribe[topic-0]].
If the job now starts but reads nothing, or no lag appears in Kafka tooling, the cause is different. See where Spark keeps its position.
A working example
Factor House publishes a PySpark job that reads Kafka, Lab 10 of Factor House Local, which packages its dependencies as a fat JAR. Factor House makes Kpow for Kafka and has no Spark product, so this page is a troubleshooting aid, not a product recommendation.
FAQ
What does Failed to find data source: kafka mean in PySpark?
Spark could not find the Kafka source on its classpath. The Kafka connector is a separate package that has to be added when the application launches.
How do I add the Kafka package to PySpark?
Pass --packages org.apache.spark:spark-sql-kafka-0-10_2.13:<spark version> to spark-submit, or to pyspark or spark-shell when experimenting, with the version of your Spark release.
Where do I find the right connector version?
In the Kafka integration guide for your Spark release. Its linking and deployment sections give the artifact name and version to pass to --packages.