【问题标题】:Flink Dynamic Update Stream jobFlink 动态更新流作业
【发布时间】:2020-06-13 15:30:38
【问题描述】:

我收到了一组关于不同主题的 Avro 格式的事件。我想使用这些并以镶木地板格式写入 s3。 我编写了一个下面的作业,它为每个事件创建一个不同的流,并从融合模式注册表中获取它的模式,以便为事件创建一个镶木地板接收器。
这工作正常,但我面临的唯一问题是每当新事件开始出现时,我必须更改 YAML 配置并每次重新启动作业。有什么方法我不必重新开始工作,它开始消耗一组新的事件。

YamlReader reader = new YamlReader(topologyConfig);
    EventTopologyConfig eventTopologyConfig = reader.read(EventTopologyConfig.class);

    long checkPointInterval = eventTopologyConfig.getCheckPointInterval();
        topics = eventTopologyConfig.getTopics();

                List<EventConfig> eventTypesList = eventTopologyConfig.getEventsType();

        CachedSchemaRegistryClient registryClient = new CachedSchemaRegistryClient(schemaRegistryUrl, 1000);


        FlinkKafkaConsumer flinkKafkaConsumer = new FlinkKafkaConsumer(topics,
        new KafkaGenericAvroDeserializationSchema(schemaRegistryUrl),
        properties);

        DataStream<GenericRecord> dataStream = streamExecutionEnvironment.addSource(flinkKafkaConsumer).name("source");

        try {
        for (EventConfig eventConfig : eventTypesList) {

        LOG.info("creating a stream for ", eventConfig.getEvent_name());

final StreamingFileSink<GenericRecord> sink = StreamingFileSink.forBulkFormat
        (path, ParquetAvroWriters.forGenericRecord(SchemaUtils.getSchema(eventConfig.getSchema_subject(), registryClient)))
        .withBucketAssigner(new EventTimeBucketAssigner())
        .build();

        DataStream<GenericRecord> outStream = dataStream.filter((FilterFunction<GenericRecord>) genericRecord -> {
        if (genericRecord != null && genericRecord.get(EVENT_NAME).toString().equals(eventConfig.getEvent_name())) {
        return true;
        }
        return false;
        });
        outStream.addSink(sink).name(eventConfig.getSink_id()).setParallelism(parallelism);

        }
        } catch (Exception e) {
        e.printStackTrace();
        }

Yaml 文件:

!com.bounce.config.EventTopologyConfig
eventsType:
  - !com.bounce.config.EventConfig
    event_name: "search_list_keyless"
    schema_subject: "search_list_keyless-com.bounce.events.keyless.bookingflow.search_list_keyless"
    topic: "search_list_keyless"

  - !com.bounce.config.EventConfig
    event_name: "bike_search_details"
    schema_subject: "bike_search_details-com.bounce.events.keyless.bookingflow.bike_search_details"
    topic: "bike_search_details"

  - !com.bounce.config.EventConfig
    event_name: "keyless_bike_lock"
    schema_subject: "analytics-keyless-com.bounce.events.keyless.bookingflow.keyless_bike_lock"
    topic: "analytics-keyless"

  - !com.bounce.config.EventConfig
      event_name: "keyless_bike_unlock"
      schema_subject: "analytics-keyless-com.bounce.events.keyless.bookingflow.keyless_bike_unlock"
      topic: "analytics-keyless"


checkPointInterval: 1200000

topics: ["search_list_keyless","bike_search_details","analytics-keyless"]

谢谢。

【问题讨论】:

    标签: java apache-flink flink-streaming


    【解决方案1】:

    我认为您想使用自定义 BucketAssigner,它使用 genericRecord.get(EVENT_NAME).toString() 值作为存储桶 ID,以及 EventTimeBucketAssigner 正在执行的任何事件时间分桶。

    那么你不需要创建多个流,它应该是动态的(每当一个新的事件名称值出现在正在写入的记录中,你就会得到一个新的输出接收器)。

    【讨论】:

    • 谢谢。由于不同的 avro 模式记录,我正在创建多个流,我必须创建具有相应模式的接收器。
    • 好的,抱歉。这让事情变得更难了。我认为您仍然可以按照我上面所说的方式进行操作,方法是创建您自己的 PartFileFactory,该 PartFileFactory 将使用您的 BucketID 类根据架构创建适当的编写器。
    • 您能否提供一些关于如何实现我自己的 PartFileFactory 的参考代码,或者您能否提供一些有用的示例代码。
    • 没有示例代码,抱歉。 Flink 源代码有一个 PartFileFactory 的实现,这是我要开始的地方。
    • 没有办法让我可以频繁地更新我的 YAML 文件。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-04-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-03-20
    相关资源
    最近更新 更多