PySpark Streaming from Cloudkarfka on Google Dataproc

Viewed 195

Service: GCP Dataproc & DataprocHub

Service: Cloudkarafka (www.cloudkarafka.com), a simple and quick way to launch your kafka service.

1: GCP-DataprocHub out of box gives Spark 2.4.8.

2: Create a Dataproc Cluster with Specific version.

gcloud shell:
gcloud dataproc clusters create dataproc-spark312  --image-version=2.0-ubuntu18  --region=us-central1 --single-node

3: Export it and save in gs bucket

gcloud dataproc clusters export dataproc-spark312 --destination dataproc-spark312.yaml --region us-central1
gsutil cp dataproc-spark312.yaml gs://gcp-learn-lib/

4: Create a ENV file

DATAPROC_CONFIGS=gs://gcp-learn-lib/dataproc-spark312.yaml
NOTEBOOKS_LOCATION=gs://gcp-learn-notebooks/notebooks
DATAPROC_LOCATIONS_LIST=a,b,c

Save as: dataproc-hub-config.env and upload to gs: bucket

5: Create a Datapro-Hub and link the above Cluster. Section: "custom env setup" key:container-env-file value: gs://gcp-learn-lib/dataproc-hub-config.env

6: Complete the Create

7: Click on Jupyter Link

8: It should show your cluster, select it. and select the region (same as the original cluster in step2)

9: Dataproc -> Cluster -> Click on Cluster (starts with hub-) -> VM Instances -> SSH Download and copy the kafka jars to

cd /usr/lib/spark/jars/
wget https://repo1.maven.org/maven2/org/apache/commons/commons-pool2/2.6.2/commons-pool2-2.6.2.jar
wget https://repo1.maven.org/maven2/org/apache/kafka/kafka-clients/2.6.0/kafka-clients-2.6.0.jar
wget https://repo1.maven.org/maven2/org/apache/spark/spark-sql-kafka-0-10_2.12/3.1.2/spark-sql-kafka-0-10_2.12-3.1.2.jar
wget https://repo1.maven.org/maven2/org/apache/spark/spark-streaming-kafka-0-10_2.12/3.1.2/spark-streaming-kafka-0-10_2.12-3.1.2.jar
wget https://repo1.maven.org/maven2/org/apache/spark/spark-streaming-kafka-0-10-assembly_2.12/3.1.2/spark-streaming-kafka-0-10-assembly_2.12-3.1.2.jar

10: Prepare JAAS and CA file

vi /tmp/cloudkarafka_gcp_oct2021.jaas
vi /tmp/cloudkarafka_gcp_oct2021.ca

11: Append to the end of the file:

sudo vi /etc/spark/conf/spark-defaults.conf
spark.executor.extraJavaOptions=-Djava.security.auth.login.config=/tmp/cloudkarafka_gcp_oct2021.jaas -Dsasl.jaas.config=/tmp/cloudkarafka_gcp_oct2021.jaas -Dssl.ca.location=/tmp/cloudkarafka_gcp_oct2021.ca
spark.driver.extraJavaOptions=-Djava.security.auth.login.config=/tmp/cloudkarafka_gcp_oct2021.jaas

12: Go back to JupyterHub. Restart Kernel - Test your code

spark.readStream.format("kafka") \
.option("kafka.bootstrap.servers", "server1:9094,serer2:9094,server3:9094") \
.option("subscribe", "foo-default") \
.option("kafka.security.protocol", "SASL_SSL") \
.option("kafka.sasl.mechanism", "SCRAM-SHA-256") \
.load() \
.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") \
.writeStream.format("kafka") \
.option("kafka.bootstrap.servers", "server1:9094,serer2:9094,server3:9094") \
.option("topic", "foo-test") \
.option("kafka.security.protocol", "SASL_SSL") \
.option("kafka.sasl.mechanism", "SCRAM-SHA-256") \
.option("checkpointLocation", "/tmp/stream/kafkatest") \
.start()

13: Enjoy

==============

0 Answers
Related