【发布时间】:2018-07-18 16:03:27
【问题描述】:
是否有人尝试使用 spark-steaming(pyspark) 作为 CDH 中 kerberos KAFKA 的消费者?
我搜索了 CDH 并找到了一些关于 Scala 的示例。
是不是说CDH不支持这个?
任何人都可以帮助解决这个问题???
【问题讨论】:
标签: pyspark apache-kafka spark-streaming kerberos cloudera-cdh
是否有人尝试使用 spark-steaming(pyspark) 作为 CDH 中 kerberos KAFKA 的消费者?
我搜索了 CDH 并找到了一些关于 Scala 的示例。
是不是说CDH不支持这个?
任何人都可以帮助解决这个问题???
【问题讨论】:
标签: pyspark apache-kafka spark-streaming kerberos cloudera-cdh
CDH 也支持基于 Pyspark 的结构化流 API 以连接受 Kerberos 保护的 Kafka 集群。即使我发现很难找到示例代码。您可以参考下面的示例代码,这些代码在 CDH 产品环境中经过良好测试和实现。
注意:以下示例代码中需要考虑的要点。
根据您的环境调整软件包版本。
在spark提交命令中提及正确的JAAS,Keytab文件位置和代码中的配置参数。
此代码已作为示例提供,用于读取启用 Kerberos 的 Kafka 集群主题并写入 HDFS 位置。
spark/bin/spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.11:2.2.0,com.databricks:spark-avro_2.11:3.2.0 --conf spark.ui.port=4055 --files /home/path/spark_jaas,/home/bdpda/bdpda.headless.keytab --conf "spark.executor.extraJavaOptions=-Djava.security.auth.login.config=/home/bdpda/spark_jaas" --conf "spark.driver.extraJavaOptions=-Djava.security.auth.login.config=/home/bdpda/spark_jaas" pysparkstructurestreaming.py
Pyspark 代码:pysparkstructurestreaming.py
from pyspark.sql import SparkSession
from pyspark.sql.functions import *
from pyspark.sql.types import *
from pyspark.streaming import StreamingContext
import time
# Spark Streaming context :
spark = SparkSession.builder.appName('PythonStreamingDirectKafkaWordCount').getOrCreate()
sc = spark.sparkContext
ssc = StreamingContext(sc, 20)
# Kafka Topic Details :
KAFKA_TOPIC_NAME_CONS = "topic_name"
KAFKA_OUTPUT_TOPIC_NAME_CONS = "topic_to_hdfs"
KAFKA_BOOTSTRAP_SERVERS_CONS = 'kafka_server:9093'
# Creating readstream DataFrame :
df = spark.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", KAFKA_BOOTSTRAP_SERVERS_CONS) \
.option("subscribe", KAFKA_TOPIC_NAME_CONS) \
.option("startingOffsets", "earliest") \
.option("kafka.security.protocol","SASL_SSL")\
.option("kafka.client.id" ,"Clinet_id")\
.option("kafka.sasl.kerberos.service.name","kafka")\
.option("kafka.ssl.truststore.location", "/home/path/kafka_trust.jks") \
.option("kafka.ssl.truststore.password", "password_rd") \
.option("kafka.sasl.kerberos.keytab","/home/path.keytab") \
.option("kafka.sasl.kerberos.principal","path") \
.load()
df1 = df.selectExpr( "CAST(value AS STRING)")
# Creating Writestream DataFrame :
df1.writeStream \
.option("path","target_directory") \
.format("csv") \
.option("checkpointLocation","chkpint_directory") \
.outputMode("append") \
.start()
ssc.awaitTermination()
【讨论】: