【问题标题】:How to identify the origin of messages in spark structured streaming with kafka as a source?如何识别以 kafka 为源的 Spark 结构化流中的消息来源?
【发布时间】:2019-06-25 03:03:51
【问题描述】:

我有一个用例,我必须在 spark 结构化流媒体中订阅 kafka 中的多个主题。然后我必须解析每条消息并从中形成一个三角洲湖表。我已经使解析器和消息(以 xml 的形式)正确解析和形成 delta-lake 表。但是,到目前为止,我只订阅了一个主题。我想订阅多个主题,并且基于主题,它应该转到专门为这个特定主题制作的解析器。所以基本上我想在所有消息处理时识别它们的主题名称,以便我可以将它们发送到所需的解析器并进一步处理。

这就是我访问来自不同主题的消息的方式。但是,我不知道如何在处理传入消息时识别它们的来源。

 val stream_dataframe = spark.readStream
  .format(ConfigSetting.getString("source"))
  .option("kafka.bootstrap.servers", ConfigSetting.getString("bootstrap_servers"))
  .option("kafka.ssl.truststore.location", ConfigSetting.getString("trustfile_location"))
  .option("kafka.ssl.truststore.password", ConfigSetting.getString("truststore_password"))
  .option("kafka.sasl.mechanism", ConfigSetting.getString("sasl_mechanism"))
  .option("kafka.security.protocol", ConfigSetting.getString("kafka_security_protocol"))
  .option("kafka.sasl.jaas.config",ConfigSetting.getString("jass_config"))
  .option("encoding",ConfigSetting.getString("encoding"))
  .option("startingOffsets",ConfigSetting.getString("starting_offset_duration"))
  .option("subscribe",ConfigSetting.getString("topics_name"))
  .option("failOnDataLoss",ConfigSetting.getString("fail_on_dataloss")) 
  .load()


 var cast_dataframe = stream_dataframe.select(col("value").cast(StringType))

 cast_dataframe =  cast_dataframe.withColumn("parsed_column",parser(col("value"))) // Parser is the udf, made to parse the xml from the topic. 

当消息在 spark 结构化流中处理时,如何识别它们的主题名称?

【问题讨论】:

    标签: apache-spark apache-kafka spark-structured-streaming


    【解决方案1】:

    根据official documentation(强调我的)

    源中的每一行都有以下架构:

    列类型


    密钥二进制
    值二进制
    主题字符串
    分区整数

    ...

    如您所见,输入主题是输出模式的一部分,无需任何特殊操作即可访问。

    【讨论】:

      猜你喜欢
      • 2017-03-31
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-04-11
      • 2019-01-29
      • 1970-01-01
      • 2021-03-14
      相关资源
      最近更新 更多