这是我的第一个卡夫卡项目(使用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
?
有人能告诉我我在这里犯了什么错误吗?
我们非常感谢您的帮助。