【问题标题】:CDH spark steaming consumer kerberos kafkaCDH火花蒸消费者kerberos kafka
【发布时间】:2018-07-18 16:03:27
【问题描述】:

是否有人尝试使用 spark-steaming(pyspark) 作为 CDH 中 kerberos KAFKA 的消费者?

我搜索了 CDH 并找到了一些关于 Scala 的示例。

是不是说CDH不支持这个?

任何人都可以帮助解决这个问题???

【问题讨论】:

    标签: pyspark apache-kafka spark-streaming kerberos cloudera-cdh


    【解决方案1】:

    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()
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-11-21
      • 2020-03-25
      • 1970-01-01
      • 2018-07-16
      • 2020-02-14
      • 2019-06-02
      • 2020-09-19
      • 2017-11-16
      相关资源
      最近更新 更多