MCPcopy Create free account

hub / github.com/coderblack/doit30_flink / functions

Functions478 in github.com/coderblack/doit30_flink

Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_09_StreamFileSinkOperator_Demo3.java:31
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flink/java/demos/ParallelismDe.java:10
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_28_ToleranceSideToSideTest.java:64
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_07_Transformation_Demos.java:29
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_09_StreamFileSinkOperator_Demo1.java:32
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo3.java:24
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_24_State_KeyedState_Demo.java:28
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_18_ChannalSelector_Partitioner_Demo.java:14
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo9_EventTimeAndWatermark3.java:32
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo12_JdbcConnectorTest1.java:10
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo22_CustomAggregateFunction.java:19
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo8_CsvFormat.java:20
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo5_CatalogDemo.java:19
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo3_TableObjectCreate.java:36
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo12_JdbcConnectorTest2.java:14
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo21_CustomScalarFunction.java:12
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo1_TableSql.java:18
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo10_KafkaConnectorDetail.java:21
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo9_EventTimeAndWatermark.java:32
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo20_Temporal_Join.java:24
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo17_TimeWindowJoin.java:23
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo6_Exercise.java:31
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo25_MetricDemos.java:14
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo14_MysqlCdcConnector.java:17
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo18_IntervalJoin.java:24
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo23_TableFunction.java:12
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo4_SqlTableCreate.java:29
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo19_ArrayJoin.java:13
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo24_TableAggregateFunction.java:61
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo7_ColumnDetail1_Sql.java:14
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo8_JsonFormat.java:20
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo13_FileSystemConnectorTest.java:16
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo11_UpsertKafkaConnectorTest2.java:15
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo2_TableApi.java:12
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo14_StreamFromToTable.java:25
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo16_TimeWindowDemo.java:26
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo11_UpsertKafkaConnectorTest.java:15
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo7_ColumnDetail2_TableApi.java:15
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo19_LookupJoin.java:23
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo24_TableAggregateFunction2.java:64
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo9_EventTimeAndWatermark2.java:34
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo18_RegularJoin.java:24
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/Exercise.java:44
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/TimerDemo.java:25
Methodmain
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/KeyedStateDemo.java:28
Methodmain
(String[] args)
kafka_course/src/test/java/RoaringBitmapTest.java:9
Methodmain
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/AdminClientDemo.java:12
Methodmain
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka自身事务机制.java:21
Methodmain
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习.java:50
Methodmain
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/ConsumerDemo.java:16
Methodmain
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/Consumer实现ExactlyOnce手段1.java:38
Methodmain
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习_消费者_判重.java:19
Methodmain
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/ProducerDemo.java:15
Methodmain
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习_消费者_Bitmap.java:18
Methodmain
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/ConsumerDemo3.java:18
Methodmain
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习_消费者.java:19
Methodmain
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/ConsumerDemo2.java:20
Methodmain
(String[] args)
datagen/src/main/java/cn/doitedu/ActionLogAutoGen.java:56
Methodmain
(String[] args)
datagen/src/main/java/cn/doitedu/ActionLogGenOne.java:22
Methodmap1
对 左流 处理的逻辑 @param value @return @throws Exception
flink_course/src/main/java/cn/doitedu/flink/java/demos/_14_StreamConnect_Union_Demo.java:56
Methodmap2
对 右流 处理的逻辑 @param value @return @throws Exception
flink_course/src/main/java/cn/doitedu/flink/java/demos/_14_StreamConnect_Union_Demo.java:68
MethodmarkIdle
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/TimestampsAndWatermarksOperator.java:174
Methodmerge
( W mergeResult, Colle
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:316
Methodmerge
(MysqlUser a, MysqlUser b)
flink_course/src/main/java/tmp/FlinkKafkaDemo.java:155
Methodmerge
批计算模式下,可能需要将多个上游的局部聚合累加器,放在下游进行全局聚合 因为需要对两个累加器进行合并 这里就是合并的逻辑 流计算模式下,不用实现! @param a An accumulator to merge @param b Another accumulator to merge @r
flink_course/src/main/java/cn/doitedu/flink/java/demos/_20_Window_Api_Demo1.java:125
Methodmerge
(Tuple2<Integer, Integer> a, Tuple2<Integer, Integer> b)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_25_State_DataStructure_Demo.java:112
Methodmerge
(MyAccumulator acc, Iterable<MyAccumulator> it)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo24_TableAggregateFunction.java:122
Methodmerge
(MyAccumulator acc, Iterable<MyAccumulator> it)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo24_TableAggregateFunction2.java:127
Methodmerge
(Tuple2<String, Integer> a, Tuple2<String, Integer> b)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/Exercise.java:198
MethodmergePartitionedState
( StateDescriptor<S, ?> stateDescriptor)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:889
MethodnotifyCheckpointAborted
(long checkpointId)
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:381
MethodnotifyCheckpointComplete
(long checkpointId)
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:376
MethodonElement
来一条数据时,需要检查watermark是否已经越过窗口结束点需要触发
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:107
MethodonEventTime
(long time)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:946
MethodonEventTime
当事件时间定时器的触发时间(窗口的结束点时间)到达了,检查是否满足触发条件 下面的方法,是定时器在调用
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:129
MethodonMerge
(TimeWindow window, OnMergeContext ctx)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:153
MethodonPartitionsAssigned
(Collection<TopicPartition> partitions)
kafka_course/src/main/java/cn/doitedu/kafka/Consumer实现ExactlyOnce手段1.java:81
MethodonPartitionsAssigned
(Collection<TopicPartition> partitions)
kafka_course/src/main/java/cn/doitedu/kafka/ConsumerDemo3.java:45
MethodonPartitionsRevoked
(Collection<TopicPartition> partitions)
kafka_course/src/main/java/cn/doitedu/kafka/Consumer实现ExactlyOnce手段1.java:74
MethodonPartitionsRevoked
(Collection<TopicPartition> partitions)
kafka_course/src/main/java/cn/doitedu/kafka/ConsumerDemo3.java:37
MethodonProcessingTime
(long timestamp)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/TimestampsAndWatermarksOperator.java:110
MethodonProcessingTime
(long time)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:942
MethodonProcessingTime
当处理时间定时器的触发时间(窗口的结束点时间)到达了,检查是否满足触发条件
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:137
MethodonTimer
定期器被触发时,会调用的方法
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/TimerDemo.java:73
Methodopen
This method is called immediately before any elements are processed, it should contain the operator's initialization logic, e.g. state initialization.
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:322
Methodopen
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/TimestampsAndWatermarksOperator.java:76
Methodopen
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:218
Methodopen
(Configuration parameters)
flink_course/src/main/java/tmp/FlinkKafkaDemo.java:177
Methodopen
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_26_State_TTL_Demo.java:85
Methodopen
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_06_CustomSourceFunction.java:178
Methodopen
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_17_ProcessFunctions_Demo.java:50
Methodopen
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_25_State_DataStructure_Demo.java:67
Methodopen
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_28_ToleranceSideToSideTest.java:150
Methodopen
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_24_State_KeyedState_Demo.java:54
Methodopen
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo25_MetricDemos.java:24
Methodopen
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/Exercise.java:107
Methodopen
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/TimerDemo.java:44
Methodopen
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/KeyedStateDemo.java:54
MethodprepareSnapshotPreBarrier
(long checkpointId)
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:335
MethodprocessBroadcastElement
(UserInfo userInfo, BroadcastProcessFunction<EventCount, UserInfo, EventUserInfo>.Context ctx, Collector<Event
flink_course/src/main/java/cn/doitedu/flink/exercise/Exercise_1.java:132
← previousnext →301–400 of 478, ranked by callers