Flink未來將與 Pulsar集成提供大規(guī)模的彈性數(shù)據(jù)處理
未來整合
Pulsar可以以不同的方式與Apache Flink集成。一些潛在的集成包括使用流式連接器為流式工作負(fù)載提供支持,并使用批量源連接器支持批量工作負(fù)載。Pulsar還提供對(duì)schema 的本地支持,可以與Flink集成并提供對(duì)數(shù)據(jù)的結(jié)構(gòu)化訪問,例如使用Flink SQL作為在Pulsar中查詢數(shù)據(jù)的方式。最后,集成這些技術(shù)的另一種方法可能包括使用Pulsar作為Flink的狀態(tài)后端。由于Pulsar具有分層架構(gòu)(Streams和Segmented Streams,由Apache Bookkeeper提供支持),因此將Pulsar用作存儲(chǔ)層并存儲(chǔ)Flink狀態(tài)變得很自然。
從體系結(jié)構(gòu)的角度來看,我們可以想象兩個(gè)框架之間的集成,它使用Apache Pulsar作為統(tǒng)一的數(shù)據(jù)層視圖,Apache Flink作為統(tǒng)一的計(jì)算和數(shù)據(jù)處理框架和API。
現(xiàn)有集成
兩個(gè)框架之間的集成正在進(jìn)行中,開發(fā)人員已經(jīng)可以通過多種方式將Pulsar與Flink結(jié)合使用。例如,Pulsar可用作Flink DataStream應(yīng)用程序中的流媒體源和流式接收器。開發(fā)人員可以將Pulsar中的數(shù)據(jù)提取到Flink作業(yè)中,該作業(yè)可以計(jì)算和處理實(shí)時(shí)數(shù)據(jù),然后將數(shù)據(jù)作為流式接收器發(fā)送回Pulsar主題。這樣的例子如下所示:
// create and configure Pulsar consumer
PulsarSourceBuilder<String>builder = PulsarSourceBuilder
.builder(new SimpleStringSchema())
.serviceUrl(serviceUrl)
.topic(inputTopic)
.subscriptionName(subscription);
SourceFunction<String> src = builder.build();
// ingest DataStream with Pulsar consumer
DataStream<String> words = env.a(chǎn)ddSource(src);
// perform computation on DataStream (here a simple WordCount)
DataStream<WordWithCount> wc = words
.flatMap((FlatMapFunction<String, WordWithCount>) (word, collector) -> {
collector.collect(new WordWithCount(word, 1));
})
.returns(WordWithCount.class)
.keyBy("word")
.timeWindow(Time.seconds(5))
.reduce((ReduceFunction<WordWithCount>) (c1, c2) ->
new WordWithCount(c1.word, c1.count + c2.count));
// emit result via Pulsar producer
wc.a(chǎn)ddSink(new FlinkPulsarProducer<>(
serviceUrl,
outputTopic,
new AuthenticationDisabled(),
wordWithCount -> wordWithCount.toString().getBytes(UTF_8),
wordWithCount -> wordWithCount.word)
);
開發(fā)人員可以利用的兩個(gè)框架之間的另一個(gè)集成包括將Pulsar用作Flink SQL或Table API查詢的流式源和流式表接收器,如下例所示:
// obtain a DataStream with words
DataStream<String> words = ...
// register DataStream as Table "words" with two attributes ("word", "ts").
// "ts" is an event-time timestamp.
tableEnvironment.registerDataStream("words", words, "word, ts.rowtime");
// create a TableSink that produces to Pulsar
TableSink sink = new PulsarJsonTableSink(
serviceUrl,
outputTopic,
new AuthenticationDisabled(),
ROUTING_KEY);
// register Pulsar TableSink as table "wc"
tableEnvironment.registerTableSink(
"wc",
sink.configure(
new String[]{"word", "cnt"},
new TypeInformation[]{Types.STRING, Types.LONG}));
// count words per 5 seconds and write result to table "wc"
tableEnvironment.sqlUpdate(
"INSERT INTO wc " +
"SELECT word, COUNT(*) AS cnt " +
"FROM words " +
"GROUP BY word, TUMBLE(ts, INTERVAL '5' SECOND)");
最后,F(xiàn)link將批量工作負(fù)載與Pulsar集成為批處理接收器,其中所有結(jié)果在Apache Flink完成靜態(tài)數(shù)據(jù)集中的計(jì)算后被推送到Pulsar。這樣的例子如下所示:
// obtain DataSet from arbitrary computation
DataSet<WordWithCount> wc = ...
// create PulsarOutputFormat instance
OutputFormat pulsarOutputFormat = new PulsarOutputFormat(
serviceUrl,
topic,
new AuthenticationDisabled(),
wordWithCount -> wordWithCount.toString().getBytes());
// write DataSet to Pulsar
wc.output(pulsarOutputFormat);
結(jié)論
Pulsar和Flink都對(duì)應(yīng)用程序的數(shù)據(jù)和計(jì)算級(jí)別如何以批量作為特殊情況流“流式傳輸”方式分享了類似的觀點(diǎn)。通過Pulsar的Segmented Streams方法和Flink在一個(gè)框架下統(tǒng)一批處理和流處理工作負(fù)載的步驟,有許多方法將這兩種技術(shù)集成在一起,以提供大規(guī)模的彈性數(shù)據(jù)處理。

發(fā)表評(píng)論
請(qǐng)輸入評(píng)論內(nèi)容...
請(qǐng)輸入評(píng)論/評(píng)論長(zhǎng)度6~500個(gè)字
圖片新聞
-
機(jī)器人奧運(yùn)會(huì)戰(zhàn)報(bào):宇樹機(jī)器人摘下首金,天工Ultra搶走首位“百米飛人”
-
存儲(chǔ)圈掐架!江波龍起訴佰維,索賠121萬(wàn)
-
長(zhǎng)安汽車母公司突然更名:從“中國(guó)長(zhǎng)安”到“辰致科技”
-
豆包前負(fù)責(zé)人喬木出軌BP后續(xù):均被辭退
-
字節(jié)AI Lab負(fù)責(zé)人李航卸任后返聘,Seed進(jìn)入調(diào)整期
-
員工持股爆雷?廣汽埃安緊急回應(yīng)
-
中國(guó)“智造”背后的「關(guān)鍵力量」
-
小米汽車研發(fā)中心重磅落地,寶馬家門口“搶人”
最新活動(dòng)更多
-
10月23日火熱報(bào)名中>> 2025是德科技創(chuàng)新技術(shù)峰會(huì)
-
10月23日立即報(bào)名>> Works With 開發(fā)者大會(huì)深圳站
-
10月24日立即參評(píng)>> 【評(píng)選】維科杯·OFweek 2025(第十屆)物聯(lián)網(wǎng)行業(yè)年度評(píng)選
-
11月27日立即報(bào)名>> 【工程師系列】汽車電子技術(shù)在線大會(huì)
-
12月18日立即報(bào)名>> 【線下會(huì)議】OFweek 2025(第十屆)物聯(lián)網(wǎng)產(chǎn)業(yè)大會(huì)
-
精彩回顧立即查看>> 【限時(shí)福利】TE 2025國(guó)際物聯(lián)網(wǎng)展·深圳站
推薦專題
- 1 先進(jìn)算力新選擇 | 2025華為算力場(chǎng)景發(fā)布會(huì)暨北京xPN伙伴大會(huì)成功舉辦
- 2 人形機(jī)器人,正狂奔在批量交付的曠野
- 3 宇樹機(jī)器人撞人事件的深度剖析:六維力傳感器如何成為人機(jī)安全的關(guān)鍵屏障
- 4 解碼特斯拉新AI芯片戰(zhàn)略 :從Dojo到AI5和AI6推理引擎
- 5 AI版“四萬(wàn)億刺激”計(jì)劃來了
- 6 2025年8月人工智能投融資觀察
- 7 8 a16z最新AI百?gòu)?qiáng)榜:硅谷頂級(jí)VC帶你讀懂全球生成式AI賽道最新趨勢(shì)
- 9 Manus跑路,大廠掉線,只能靠DeepSeek了
- 10 地平線的野心:1000萬(wàn)套HSD上車