【问题标题】:Add configuration parameters - spark & Kafka : acks and compression添加配置参数 - spark & Kafka:acks 和压缩
【发布时间】:2019-06-21 19:47:05
【问题描述】:

我想向我的应用程序 spark 和 Kafka 添加一些参数,以便将 Dataframe 写入主题 kafka。

我没有在 spark-kafka 文档中找到 acks 和 compression.codec

   .write
   .format("kafka")
   .option("kafka.sasl.mechanism", Config.KAFKA_SASL_MECHANISM)
   .option("kafka.security.protocol", Config.KAFKA_SECURITY_PROTOCOL)
   .option("kafka.sasl.jaas.config", KAFKA_JAAS_CONFIG)
   .option("kafka.bootstrap.servers", KAFKA_BOOTSTRAP)
   .option("fetchOffset.numRetries", 6)
   .option("acks","all")
   .option("compression.codec","lz4")
   .option("kafka.request.timeout.ms", 120000)
   .option("topic", topic)
   .save()```

【问题讨论】:

    标签: scala apache-spark apache-kafka


    【解决方案1】:

    您可以使用此特定属性来定义您的序列化程序: default.value.serde

    【讨论】:

    • 消息密钥呢?我没有在 Spark 文档中看到这个属性
    【解决方案2】:

    对于序列化程序,创建一个案例类或一到三列数据框,其中仅包含 keyvalueArray[Byte] 字段(字符串也可以使用)。然后topic 字符串字段。如果你只需要Kafka值,那么你只需要一列Dataframe

    在写入 Kafka 之前,您需要映射当前数据以将其全部序列化。

    然后,文档确实说任何其他生产者属性都只是以 kafka. 为前缀

    更多信息在这里https://spark.apache.org/docs/latest/structured-streaming-kafka-integration.html#writing-data-to-kafka

    对于 SASL 属性,我认为您需要在提交期间使用 spark.executor.options 并使用 --files 传递关键选项卡或 jaas 文件,尽管

    【讨论】:

    • 我想通过 .option("","") 添加 Acks 属性
    • kafka.acks 应该是它,根据文档
    猜你喜欢
    • 2016-07-16
    • 1970-01-01
    • 2011-11-30
    • 2016-09-20
    • 2014-03-20
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多