Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/coderblack/doit30_flink
/ functions
Functions
478 in github.com/coderblack/doit30_flink
⨍
Functions
478
◇
Types & classes
167
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_09_StreamFileSinkOperator_Demo3.java:31
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flink/java/demos/ParallelismDe.java:10
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_28_ToleranceSideToSideTest.java:64
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_07_Transformation_Demos.java:29
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_09_StreamFileSinkOperator_Demo1.java:32
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo3.java:24
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_24_State_KeyedState_Demo.java:28
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_18_ChannalSelector_Partitioner_Demo.java:14
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo9_EventTimeAndWatermark3.java:32
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo12_JdbcConnectorTest1.java:10
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo22_CustomAggregateFunction.java:19
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo8_CsvFormat.java:20
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo5_CatalogDemo.java:19
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo3_TableObjectCreate.java:36
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo12_JdbcConnectorTest2.java:14
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo21_CustomScalarFunction.java:12
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo1_TableSql.java:18
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo10_KafkaConnectorDetail.java:21
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo9_EventTimeAndWatermark.java:32
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo20_Temporal_Join.java:24
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo17_TimeWindowJoin.java:23
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo6_Exercise.java:31
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo25_MetricDemos.java:14
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo14_MysqlCdcConnector.java:17
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo18_IntervalJoin.java:24
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo23_TableFunction.java:12
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo4_SqlTableCreate.java:29
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo19_ArrayJoin.java:13
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo24_TableAggregateFunction.java:61
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo7_ColumnDetail1_Sql.java:14
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo8_JsonFormat.java:20
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo13_FileSystemConnectorTest.java:16
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo11_UpsertKafkaConnectorTest2.java:15
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo2_TableApi.java:12
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo14_StreamFromToTable.java:25
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo16_TimeWindowDemo.java:26
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo11_UpsertKafkaConnectorTest.java:15
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo7_ColumnDetail2_TableApi.java:15
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo19_LookupJoin.java:23
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo24_TableAggregateFunction2.java:64
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo9_EventTimeAndWatermark2.java:34
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo18_RegularJoin.java:24
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/Exercise.java:44
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/TimerDemo.java:25
Method
main
(String[] args)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/KeyedStateDemo.java:28
Method
main
(String[] args)
kafka_course/src/test/java/RoaringBitmapTest.java:9
Method
main
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/AdminClientDemo.java:12
Method
main
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka自身事务机制.java:21
Method
main
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习.java:50
Method
main
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/ConsumerDemo.java:16
Method
main
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/Consumer实现ExactlyOnce手段1.java:38
Method
main
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习_消费者_判重.java:19
Method
main
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/ProducerDemo.java:15
Method
main
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习_消费者_Bitmap.java:18
Method
main
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/ConsumerDemo3.java:18
Method
main
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习_消费者.java:19
Method
main
(String[] args)
kafka_course/src/main/java/cn/doitedu/kafka/ConsumerDemo2.java:20
Method
main
(String[] args)
datagen/src/main/java/cn/doitedu/ActionLogAutoGen.java:56
Method
main
(String[] args)
datagen/src/main/java/cn/doitedu/ActionLogGenOne.java:22
Method
map1
对 左流 处理的逻辑 @param value @return @throws Exception
flink_course/src/main/java/cn/doitedu/flink/java/demos/_14_StreamConnect_Union_Demo.java:56
Method
map2
对 右流 处理的逻辑 @param value @return @throws Exception
flink_course/src/main/java/cn/doitedu/flink/java/demos/_14_StreamConnect_Union_Demo.java:68
Method
markIdle
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/TimestampsAndWatermarksOperator.java:174
Method
merge
( W mergeResult, Colle
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:316
Method
merge
(MysqlUser a, MysqlUser b)
flink_course/src/main/java/tmp/FlinkKafkaDemo.java:155
Method
merge
批计算模式下,可能需要将多个上游的局部聚合累加器,放在下游进行全局聚合 因为需要对两个累加器进行合并 这里就是合并的逻辑 流计算模式下,不用实现! @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
Method
merge
(Tuple2<Integer, Integer> a, Tuple2<Integer, Integer> b)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_25_State_DataStructure_Demo.java:112
Method
merge
(MyAccumulator acc, Iterable<MyAccumulator> it)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo24_TableAggregateFunction.java:122
Method
merge
(MyAccumulator acc, Iterable<MyAccumulator> it)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo24_TableAggregateFunction2.java:127
Method
merge
(Tuple2<String, Integer> a, Tuple2<String, Integer> b)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/Exercise.java:198
Method
mergePartitionedState
( StateDescriptor<S, ?> stateDescriptor)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:889
Method
notifyCheckpointAborted
(long checkpointId)
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:381
Method
notifyCheckpointComplete
(long checkpointId)
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:376
Method
onElement
来一条数据时,需要检查watermark是否已经越过窗口结束点需要触发
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:107
Method
onEventTime
(long time)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:946
Method
onEventTime
当事件时间定时器的触发时间(窗口的结束点时间)到达了,检查是否满足触发条件 下面的方法,是定时器在调用
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:129
Method
onMerge
(TimeWindow window, OnMergeContext ctx)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:153
Method
onPartitionsAssigned
(Collection<TopicPartition> partitions)
kafka_course/src/main/java/cn/doitedu/kafka/Consumer实现ExactlyOnce手段1.java:81
Method
onPartitionsAssigned
(Collection<TopicPartition> partitions)
kafka_course/src/main/java/cn/doitedu/kafka/ConsumerDemo3.java:45
Method
onPartitionsRevoked
(Collection<TopicPartition> partitions)
kafka_course/src/main/java/cn/doitedu/kafka/Consumer实现ExactlyOnce手段1.java:74
Method
onPartitionsRevoked
(Collection<TopicPartition> partitions)
kafka_course/src/main/java/cn/doitedu/kafka/ConsumerDemo3.java:37
Method
onProcessingTime
(long timestamp)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/TimestampsAndWatermarksOperator.java:110
Method
onProcessingTime
(long time)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:942
Method
onProcessingTime
当处理时间定时器的触发时间(窗口的结束点时间)到达了,检查是否满足触发条件
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:137
Method
onTimer
定期器被触发时,会调用的方法
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/TimerDemo.java:73
Method
open
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
Method
open
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/TimestampsAndWatermarksOperator.java:76
Method
open
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:218
Method
open
(Configuration parameters)
flink_course/src/main/java/tmp/FlinkKafkaDemo.java:177
Method
open
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_26_State_TTL_Demo.java:85
Method
open
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_06_CustomSourceFunction.java:178
Method
open
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_17_ProcessFunctions_Demo.java:50
Method
open
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_25_State_DataStructure_Demo.java:67
Method
open
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_28_ToleranceSideToSideTest.java:150
Method
open
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_24_State_KeyedState_Demo.java:54
Method
open
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo25_MetricDemos.java:24
Method
open
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/Exercise.java:107
Method
open
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/TimerDemo.java:44
Method
open
(Configuration parameters)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/KeyedStateDemo.java:54
Method
prepareSnapshotPreBarrier
(long checkpointId)
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:335
Method
processBroadcastElement
(UserInfo userInfo, BroadcastProcessFunction<EventCount, UserInfo, EventUserInfo>.Context ctx, Collector<Event
flink_course/src/main/java/cn/doitedu/flink/exercise/Exercise_1.java:132
← previous
next →
301–400 of 478, ranked by callers