2020-08-28

import org.apache.commons.lang3.StringUtils;
import org.apache.flink.api.common.typeinfo.Types;
import org.apache.flink.runtime.state.CheckpointListener;
import org.apache.flink.runtime.state.FunctionInitializationContext;
import org.apache.flink.runtime.state.FunctionSnapshotContext;
import org.apache.flink.streaming.api.checkpoint.CheckpointedFunction;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.ProcessFunction;
import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
import org.apache.flink.streaming.api.functions.source.RichParallelSourceFunction;
import org.apache.flink.streaming.api.operators.AbstractStreamOperator;
import org.apache.flink.streaming.api.operators.OneInputStreamOperator;
import org.apache.flink.streaming.runtime.streamrecord.StreamRecord;
import org.apache.flink.util.Collector;

import java.text.SimpleDateFormat;
import java.util.ArrayList;
import java.util.Date;
import java.util.List;

public class FlinkCheckpointTest {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment steamEnv = StreamExecutionEnvironment.getExecutionEnvironment();
steamEnv.enableCheckpointing(1000L*2);
steamEnv
.addSource(new FSource()).setParallelism(4)
.transform("開始事務(wù)", Types.STRING,new FStart()).setParallelism(1)
.process(new FCombine()).name("事務(wù)預(yù)處理").setParallelism(4)
.addSink(new FSubmit()).name("提交事務(wù)").setParallelism(1)
;
steamEnv.execute("test");
}

static class FSource extends RichParallelSourceFunction<String>{
@Override
public void run(SourceContext<String> sourceContext) throws Exception {
int I =0;
while (true){
I = I + 1;
sourceContext.collect(Thread.currentThread().getId() +"-" +I);
Thread.sleep(1000);
}
}
@Override
public void cancel() {}
}

static class FStart extends AbstractStreamOperator<String> implements OneInputStreamOperator<String,String>{
   volatile Long ckid = 0L;
    @Override
    public void processElement(StreamRecord<String> streamRecord) throws Exception {
        String value = streamRecord.getValue() + "-" + ckid;
        streamRecord.replace(value);
        log("收到數(shù)據(jù): " + value + "..ckid:" + ckid);
        output.collect(streamRecord);
    }
    @Override
    public void prepareSnapshotPreBarrier(long checkpointId) throws Exception {
        log("開啟事務(wù): " + checkpointId);
        ckid = checkpointId;
        super.prepareSnapshotPreBarrier(checkpointId);
    }
}

static class FCombine extends ProcessFunction<String,String> implements CheckpointedFunction {
    List ls = new ArrayList<String>();
    Collector<String> collector =null;
    volatile Long ckid = 0L;

    @Override
    public void snapshotState(FunctionSnapshotContext functionSnapshotContext) throws Exception {
        StringBuffer sb = new StringBuffer();
        ls.forEach(x->{sb.append(x).append(";");});
        log("批處理 " + functionSnapshotContext.getCheckpointId() + ": 時(shí)收到數(shù)據(jù):" + sb.toString());
        Thread.sleep(5*1000);
        collector.collect("["+sb.toString() + "-"+ckid+"]");
        ls.clear();
        log("批處理 " + functionSnapshotContext.getCheckpointId() + " 完成");
        //Thread.sleep(5*1000);
        //Thread.sleep(20*1000);
    }
    @Override
    public void initializeState(FunctionInitializationContext functionInitializationContext) throws Exception {        }
    @Override
    public void processElement(String s, Context context, Collector<String> out) throws Exception {
        if(StringUtils.isNotBlank(s)){
            ls.add(s);
        }
        log("收到數(shù)據(jù) :" + s + "; 這批數(shù)據(jù)大小為:" + ls.size() + "..ckid:" + ckid);
        if(collector ==null){
            collector = out;
        }
    }
}

static class FSubmit extends RichSinkFunction<String> implements  CheckpointedFunction, CheckpointListener {
    List ls = new ArrayList<String>();
    volatile Long ckid = 0L;
    @Override
    public void notifyCheckpointComplete(long l) throws Exception {
        ckid = l;
        StringBuffer sb = new StringBuffer();
        ls.forEach(x->{sb.append(x).append("||");});
        log("submit notifyCheckpointComplete " + l + " over data:list size" + ls.size()+ "; detail" + sb.toString());
        Thread.sleep(100000);
        ls.clear();
    }
    @Override
    public void invoke(String value, Context context) throws Exception {
        if(StringUtils.isNotBlank(value)){
            ls.add(value);
        }
        log("收到數(shù)據(jù) :" + value + " list zie:" + ls.size() + "..ckid:" + ckid);
    }

    @Override
    public void snapshotState(FunctionSnapshotContext context) throws Exception {
        StringBuffer sb = new StringBuffer();
        ls.forEach(x->{sb.append(x).append("||");});
        log("submit snapshotState " + context.getCheckpointId() + " over data:list size" + ls.size()+ "; detail" + sb.toString());
        ls.clear();
    }

    @Override
    public void initializeState(FunctionInitializationContext context) throws Exception {

    }
}
public static void log(String s){
    String name = Thread.currentThread().getName();
    System.out.println(new SimpleDateFormat("HH:mm:ss.SSS").format(new Date())+":"+name + ":" + s);
}

}

最后編輯于
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請(qǐng)聯(lián)系作者
  • 序言:七十年代末欺嗤,一起剝皮案震驚了整個(gè)濱河市即纲,隨后出現(xiàn)的幾起案子,更是在濱河造成了極大的恐慌拼坎,老刑警劉巖圣拄,帶你破解...
    沈念sama閱讀 211,743評(píng)論 6 492
  • 序言:濱河連續(xù)發(fā)生了三起死亡事件荣刑,死亡現(xiàn)場(chǎng)離奇詭異措嵌,居然都是意外死亡,警方通過查閱死者的電腦和手機(jī)谤辜,發(fā)現(xiàn)死者居然都...
    沈念sama閱讀 90,296評(píng)論 3 385
  • 文/潘曉璐 我一進(jìn)店門蓄坏,熙熙樓的掌柜王于貴愁眉苦臉地迎上來价捧,“玉大人,你說我怎么就攤上這事涡戳〗狍” “怎么了?”我有些...
    開封第一講書人閱讀 157,285評(píng)論 0 348
  • 文/不壞的土叔 我叫張陵渔彰,是天一觀的道長(zhǎng)嵌屎。 經(jīng)常有香客問我,道長(zhǎng)恍涂,這世上最難降的妖魔是什么编整? 我笑而不...
    開封第一講書人閱讀 56,485評(píng)論 1 283
  • 正文 為了忘掉前任,我火速辦了婚禮乳丰,結(jié)果婚禮上,老公的妹妹穿的比我還像新娘内贮。我一直安慰自己产园,他們只是感情好,可當(dāng)我...
    茶點(diǎn)故事閱讀 65,581評(píng)論 6 386
  • 文/花漫 我一把揭開白布夜郁。 她就那樣靜靜地躺著什燕,像睡著了一般。 火紅的嫁衣襯著肌膚如雪竞端。 梳的紋絲不亂的頭發(fā)上屎即,一...
    開封第一講書人閱讀 49,821評(píng)論 1 290
  • 那天,我揣著相機(jī)與錄音事富,去河邊找鬼技俐。 笑死,一個(gè)胖子當(dāng)著我的面吹牛统台,可吹牛的內(nèi)容都是我干的雕擂。 我是一名探鬼主播,決...
    沈念sama閱讀 38,960評(píng)論 3 408
  • 文/蒼蘭香墨 我猛地睜開眼贱勃,長(zhǎng)吁一口氣:“原來是場(chǎng)噩夢(mèng)啊……” “哼井赌!你這毒婦竟也來了?” 一聲冷哼從身側(cè)響起贵扰,我...
    開封第一講書人閱讀 37,719評(píng)論 0 266
  • 序言:老撾萬榮一對(duì)情侶失蹤仇穗,失蹤者是張志新(化名)和其女友劉穎,沒想到半個(gè)月后戚绕,有當(dāng)?shù)厝嗽跇淞掷锇l(fā)現(xiàn)了一具尸體纹坐,經(jīng)...
    沈念sama閱讀 44,186評(píng)論 1 303
  • 正文 獨(dú)居荒郊野嶺守林人離奇死亡,尸身上長(zhǎng)有42處帶血的膿包…… 初始之章·張勛 以下內(nèi)容為張勛視角 年9月15日...
    茶點(diǎn)故事閱讀 36,516評(píng)論 2 327
  • 正文 我和宋清朗相戀三年舞丛,在試婚紗的時(shí)候發(fā)現(xiàn)自己被綠了恰画。 大學(xué)時(shí)的朋友給我發(fā)了我未婚夫和他白月光在一起吃飯的照片宾茂。...
    茶點(diǎn)故事閱讀 38,650評(píng)論 1 340
  • 序言:一個(gè)原本活蹦亂跳的男人離奇死亡,死狀恐怖拴还,靈堂內(nèi)的尸體忽然破棺而出跨晴,到底是詐尸還是另有隱情,我是刑警寧澤片林,帶...
    沈念sama閱讀 34,329評(píng)論 4 330
  • 正文 年R本政府宣布端盆,位于F島的核電站,受9級(jí)特大地震影響费封,放射性物質(zhì)發(fā)生泄漏焕妙。R本人自食惡果不足惜,卻給世界環(huán)境...
    茶點(diǎn)故事閱讀 39,936評(píng)論 3 313
  • 文/蒙蒙 一弓摘、第九天 我趴在偏房一處隱蔽的房頂上張望焚鹊。 院中可真熱鬧,春花似錦韧献、人聲如沸末患。這莊子的主人今日做“春日...
    開封第一講書人閱讀 30,757評(píng)論 0 21
  • 文/蒼蘭香墨 我抬頭看了看天上的太陽璧针。三九已至,卻和暖如春渊啰,著一層夾襖步出監(jiān)牢的瞬間探橱,已是汗流浹背。 一陣腳步聲響...
    開封第一講書人閱讀 31,991評(píng)論 1 266
  • 我被黑心中介騙來泰國打工绘证, 沒想到剛下飛機(jī)就差點(diǎn)兒被人妖公主榨干…… 1. 我叫王不留隧膏,地道東北人。 一個(gè)月前我還...
    沈念sama閱讀 46,370評(píng)論 2 360
  • 正文 我出身青樓嚷那,卻偏偏與公主長(zhǎng)得像私植,于是被迫代替她去往敵國和親。 傳聞我的和親對(duì)象是個(gè)殘疾皇子车酣,可洞房花燭夜當(dāng)晚...
    茶點(diǎn)故事閱讀 43,527評(píng)論 2 349