• pyflink读取kafka数据写入mysql实例


    依赖包下载

    https://repo.maven.apache.org/maven2/org/apache/flink/flink-sql-connector-kafka/1.17.1/

    版本

    flink:1.16.0

    kafka:2.13-3.2.0

    实例

    1. import logging
    2. import sys
    3. from pyflink.common import Types
    4. from pyflink.datastream import StreamExecutionEnvironment
    5. from pyflink.datastream.connectors.jdbc import JdbcSink, JdbcConnectionOptions
    6. from pyflink.datastream.connectors.kafka import FlinkKafkaProducer, FlinkKafkaConsumer
    7. from pyflink.datastream.formats.json import JsonRowSerializationSchema, JsonRowDeserializationSchema
    8. def write_to_kafka(env):
    9. type_info = Types.ROW([Types.INT(), Types.STRING()])
    10. ds = env.from_collection(
    11. [(1, 'hi'), (2, 'hello'), (3, 'hi'), (4, 'hello'), (5, 'hi'), (6, 'hello'), (6, 'hello')],
    12. type_info=type_info)
    13. serialization_schema = JsonRowSerializationSchema.Builder() \
    14. .with_type_info(type_info) \
    15. .build()
    16. kafka_producer = FlinkKafkaProducer(
    17. topic='test_json_topic',
    18. serialization_schema=serialization_schema,
    19. producer_config={'security.protocol': 'SASL_PLAINTEXT', 'sasl.mechanism': 'PLAIN', 'bootstrap.servers': '192.168.1.110:9092', 'group.id': 'test-consumer-group', 'sasl.jaas.config': 'org.apache.kafka.common.security.scram.ScramLoginModule required username=\"aaaaaaaaa\" password=\"bbbbbbb\";'}
    20. )
    21. # note that the output type of ds must be RowTypeInfo
    22. ds.add_sink(kafka_producer)
    23. env.execute()
    24. def read_from_kafka(env):
    25. deserialization_schema = JsonRowDeserializationSchema.Builder() \
    26. .type_info(Types.ROW([Types.INT(), Types.STRING()])) \
    27. .build()
    28. kafka_consumer = FlinkKafkaConsumer(
    29. topics='test_json_topic',
    30. deserialization_schema=deserialization_schema,
    31. properties={'security.protocol': 'SASL_PLAINTEXT', 'sasl.mechanism': 'PLAIN', 'bootstrap.servers': '192.168.1.110:9092', 'group.id': 'test-consumer-group', 'sasl.jaas.config': 'org.apache.kafka.common.security.scram.ScramLoginModule required username=\"aaaaa\" password=\"bbbbbb\";'}
    32. )
    33. kafka_consumer.set_start_from_earliest()
    34. env.add_source(kafka_consumer).print()
    35. env.execute()
    36. def wirte_data_todb(env, data):
    37. type_info = Types.ROW([Types.INT(), Types.STRING()])
    38. env.from_collection(
    39. [(101, "Stream Processing with Apache Flink"),
    40. (102, "Streaming Systems"),
    41. (103, "Designing Data-Intensive Applications"),
    42. (104, "Kafka: The Definitive Guide")
    43. ], type_info=type_info) \
    44. .add_sink(
    45. JdbcSink.sink(
    46. "insert into flink (id, title) values (?, ?)",
    47. type_info,
    48. JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
    49. .with_url('jdbc:mysql://192.168.1.110:23006/test')
    50. .with_driver_name('com.mysql.jdbc.Driver')
    51. .with_user_name('sino')
    52. .with_password('Caib@sgcc-56')
    53. .build()
    54. ))
    55. env.execute()
    56. if __name__ == '__main__':
    57. logging.basicConfig(stream=sys.stdout, level=logging.INFO, format="%(message)s")
    58. env = StreamExecutionEnvironment.get_execution_environment()
    59. #env.add_jars("file:///opt/flink/flink-sql-connector-kafka-1.15.0.jar")
    60. #env.add_jars("file:///opt/flink/kafka-clients-2.8.1.jar")
    61. #env.add_jars("file:///opt/flink/flink-connector-jdbc-1.16.0.jar")
    62. #env.add_jars("file:///opt/flink/mysql-connector-java-8.0.29.jar")
    63. print("start reading data from kafka")
    64. read_from_kafka(env)
    65. #wirte_data_todb(env, "")
  • 相关阅读:
    Maven编程环境搭建以及VS code Maven设置
    petalinux_zynq7 驱动DAC以及ADC模块之五:nodejs+vue3实现web网页波形显示
    数据结构-单链表操作
    mfc入门基础(六)创建模态对话框与非模态对话框
    二,几何相交-5,BO算法实现--(3)事件和操作
    Vuepress 三分钟搭建一个精美的文档或博客
    【目标检测】YOLOv5遇上知识蒸馏
    Visual Studio Code配置C/C++开发环境
    javascript的call、apply、bind的实现
    电力社区电力故障,潜在风险如何避免?
  • 原文地址:https://blog.csdn.net/u012206617/article/details/133700990