abhisai muramsetti
06/10/2024, 6:11 AMenv = StreamExecutionEnvironment.get_execution_environment()
table_env = StreamTableEnvironment.create(env)
env.add_jars("file:///home/muramsettisai/miniconda3/envs/new_flink_env/lib/python3.10/site-packages/pyflink/flink-sql-connector-kafka-3.1.0-1.18.jar") # Replace with your actual path
kafka_source1 = KafkaSource.builder() \
.set_bootstrap_servers("-----") \
.set_topics("----") \
.set_group_id('----')\
.set_value_only_deserializer(SimpleStringSchema()).set_starting_offsets(KafkaOffsetsInitializer.committed_offsets(KafkaOffsetResetStrategy.LATEST)) \
.build()
data_stream1 = env.from_source(
source=kafka_source1,
watermark_strategy=WatermarkStrategy.no_watermarks(),
source_name="Kafka Source 1",
)Job Thomas
06/10/2024, 8:30 AMabhisai muramsetti
06/10/2024, 9:13 AMabhisai muramsetti
06/10/2024, 9:14 AMabhisai muramsetti
06/10/2024, 9:17 AMPedro Mázala
06/10/2024, 10:38 AMlatest instead of earliestabhisai muramsetti
06/10/2024, 10:45 AMi mentioned latest in .set_value_only_deserializer(SimpleStringSchema()).set_starting_offsets(KafkaOffsetsInitializer.committed_offsets(KafkaOffsetResetStrategy.LATEST))