亚洲在线久爱草,狠狠天天香蕉网,天天搞日日干久草,伊人亚洲日本欧美

為了賬號安全,請及時綁定郵箱和手機立即綁定
已解決430363個問題,去搜搜看,總會有你想問的

應用于每個輸入

應用于每個輸入

楊__羊羊 2021-10-13 10:59:01
我FlinkKafkaConsumer011訂閱了一個主題。我希望apply在每個 kafka 消費者消息上處理 ( ),因此自定義在每個元素FooTrigger上返回TriggerResult.FIRE。以下代碼有效,我只是對timeWindowAll(Time.minutes(1)). 看起來我做錯了什么。// set up streaming execution environmentStreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();env.setStreamTimeCharacteristic(TimeCharacteristic.IngestionTime);// create a Kafka consumerFlinkKafkaConsumer011<Foo> consumer =  new FlinkKafkaConsumer011<>(    "topic",    new Foo.FooSchema(),    props);   // Properties object// create Kafka consumer data sourceDataStream<FooTuple> trades = env.addSource(consumer)    .timeWindowAll(Time.minutes(1))    .trigger(new FooTrigger())    .evictor(new FooEvictor())    .apply(new CreateFoos());
查看完整描述

1 回答

?
當年話下

TA貢獻1890條經驗 獲得超9個贊

如果您的目標是將函數應用于流中的每個事件,ProcessFunction那么在 Flink 中使用a將是一種更自然的方法?;蛘咴诟唵蔚那闆r下,您可以使用地圖或平面地圖,或其豐富的變體,即 RichMapFunction 或 RichFlatMapFunction —— 這完全取決于您要嘗試做什么。

使用 map 或 flatmap,您可以執行無狀態的一對一或一對多轉換,它們的豐富變體可以使用鍵控狀態,而 ProcessFunction 可以使用狀態和計時器(前提是流已被鍵控)。

timeWindowAll 適用于流未按鍵分區的情況,并且您希望按持續時間定義的批處理進行非并行處理(對于鍵控并行窗口,請改用 timeWindow)。如果您只想在數據到達時對其進行處理,那么窗口化會增加不必要的復雜性。


查看完整回答
反對 回復 2021-10-13
  • 1 回答
  • 0 關注
  • 95 瀏覽
慕課專欄
更多

添加回答

舉報

0/150
提交
取消
微信客服

購課補貼
聯系客服咨詢優惠詳情

幫助反饋 APP下載

慕課網APP
您的移動學習伙伴

公眾號

掃描二維碼
關注慕課網微信公眾號