代码之家  ›  专栏  ›  技术社区  ›  Metadata

为什么kafka在主题被读取一次后仍返回相同的数据?

  •  0
  • Metadata  · 技术社区  · 4 年前

    这是我的第一个卡夫卡项目(使用Spark流媒体) 我正在尝试阅读一个卡夫卡主题,该主题是从上游来源获取数据。

    他们正在以以下方式将数据推入卡夫卡主题:

    def kafka_ingest(df: DataFrame, kafkaconfig: dict, topic_name: str):
        jaas_config = kafkaconfig['jaas_config'] + \
        f" oauth.client.id='{kafkaconfig['client_id']}'" + \
        f" oauth.client.secret='{kafkaconfig['client_secret']}'" + \
        f" oauth.token.endpoint.uri='{kafkaconfig['endpoint_uri']}'" + \
        " oauth.max.token.expiry.seconds='30000' ;"
        df.write.format('kafka') \
          .option('kafka.bootstrap.servers', kafkaconfig['kafka_broker']) \
          .option('kafka.batch.size', kafkaconfig['kafka_batch_size']) \
          .option('retries', kafkaconfig['retries']) \
          .option('kafka.max.block.ms', kafkaconfig['kafka_max_block_ms']) \
          .option('kafka.metadata.max.age.ms', kafkaconfig['kafka_metadata_max_age_ms']) \
          .option('kafka.request.timeout.ms', kafkaconfig['kafka_request_timeout_ms']) \
          .option('kafka.linger.ms', kafkaconfig['kafka_linger_ms']) \
          .option('kafka.delivery.timeout.ms', kafkaconfig['kafka_delivery_timeout_ms']) \
          .option('acks', kafkaconfig['acks']) \
          .option('kafka.security.protocol', kafkaconfig['kafka_security_protocol']) \
          .option('kafka.sasl.jaas.config', jaas_config) \
          .option('kafka.sasl.login.callback.handler.class', kafkaconfig['kafka_sasl_login_callback_handler_class']) \
          .option('kafka.sasl.mechanism', kafkaconfig['kafka_sasl_mechanism']) \
          .option('topic', topic_name) \
          .save()
    

    当将数据摄取到Kafka中时,我在foreachBatch方法中使用上述方法,其中我还提到了下面给出的相应检查点。

    def write_stream_batches(kafka_df: DataFrame, checkpoint_location: str):
        kafka_df.writeStream \
          .format('kafka') \
          .foreachBatch(join_kafka_streams_po_denorm) \
          .option('checkpointLocation', checkpoint_location) \
          .start() \
          .awaitTermination()
    
    def join_kafka_streams_po_denorm(kafka_df: DataFrame, batch_id: int):
        final_df = kafka_df.some_transformations
        kafka_ingest(kafka_ingest, kafkaconfig, topic_name)
    

    我阅读该主题的数据如下:

    def extract_kafka_data(kafka_config: dict, topic_name: str, column_schema: str, checkpoint_location: str):
        schema = extract_schema(column_schema)
        jass_config = kafka_config['jaas_config'] \
                      + " oauth.token.endpoint.uri=" + '"' + kafka_config['endpoint_uri'] + '"' \
                      + " oauth.client.id=" + '"' + kafka_config['client_id'] + '"' \
                      + " oauth.client.secret=" + '"' + kafka_config['client_secret'] + '" ;'
        stream_df = spark.readStream \
            .format('kafka') \
            .option('kafka.bootstrap.servers', kafka_config['kafka_broker']) \
            .option('subscribe', topic_name) \
            .option('kafka.security.protocol', kafka_config['kafka_security_protocol']) \
            .option('kafka.sasl.mechanism', kafka_config['kafka_sasl_mechanism']) \
            .option('kafka.sasl.jaas.config', jass_config) \
            .option('kafka.sasl.login.callback.handler.class', kafka_config['kafka_sasl_login_callback_handler_class']) \
            .option('startingOffsets', 'earliest') \
            .option('fetchOffset.retryIntervalMs', kafka_config['kafka_fetch_offset_retry_intervalms']) \
            .option('fetchOffset.numRetries', kafka_config['retries']) \
            .option('failOnDataLoss', 'False') \
            .option('checkpointLocation', checkpoint_location) \
            .load() \
            .select(from_json(col('value').cast('string'), schema).alias("json_dta")).selectExpr('json_dta.*')
        return stream_df
    

    每次我显示数据帧中的数据时,我都会看到相同的数据返回: 阅读1:

    df = extract_kafka_data(kafka_config, topic_name, column_schema, checkpoint_location)
    display(df)
    

    输出

    +---------+-------+
    |dept_name|dept_id|
    +---------+-------+
    |Finance  |10     |
    |Marketing|20     |
    |Sales    |30     |
    |IT       |40     |
    +---------+-------+
    

    我在我的主题中有4条记录,我使用该方法摄入了这些记录: kafka_ingest 。现在我已经阅读了所有4条记录,如果我再次阅读该主题,我预计不会有输出。 阅读2:

    df=extract_kafka_data(kafka_config,topic_name,column_schema,checkpoint_location)
    显示器(df)
    

    输出

    +---------+-------+
    |dept_name | dept_id|
    +---------+-------+
    |财务|10|
    |市场营销|20|
    |销售额|30|
    |它|40|
    +---------+-------+
    

    一旦我读取了Read1中的数据,根据主题中的偏移机制,我就不应该再看到相同的数据了。 但该作业返回的数据与以前相同 Read 1 .

    我设置的方式有问题吗 offset strategy 以及的用法 checkpointing ? 有人能告诉我我在这里犯了什么错误吗? 我们非常感谢您的帮助。

    0 回复  |  直到 4 年前