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
↓ 1 callers
Method
deleteCleanupTimer
Deletes the cleanup timer set for the contents of the provided window. @param window the window whose state to discard
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:652
↓ 1 callers
Method
extractTimestamp
(UserSlotGame ele, long l)
flink_course/src/main/java/tmp/FlinkKafkaDemo.java:89
↓ 1 callers
Method
flatMap
(String s, Collector<Tuple2<String, Integer>> collector)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_01_StreamWordCount.java:49
↓ 1 callers
Method
flatMap
(String value, Collector<Tuple2<String, Integer>> out)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_03_StreamBatchWordCount.java:34
↓ 1 callers
Method
flatMap
(UserInfo value, Collector<UserFriendInfo> out)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_07_Transformation_Demos.java:72
↓ 1 callers
Method
flatMap
(String value, Collector<String> out)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_18_ChannalSelector_Partitioner_Demo.java:27
↓ 1 callers
Method
flatMap
(String value, Collector<String> out)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/KeyedStateDemo.java:35
↓ 1 callers
Method
genBatchToConsole
(List<LogBeanWrapper> wrapperedUsers, int threads , Collector collector)
datagen/src/main/java/cn/doitedu/ActionLogAutoGen.java:81
↓ 1 callers
Method
genData
()
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习.java:76
↓ 1 callers
Method
getAggregatingState
( AggregatingStateDescriptor<IN, ACC, OUT> stateProperties)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:731
↓ 1 callers
Method
getContainingTask
()
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:404
↓ 1 callers
Method
getGame_id
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:47
↓ 1 callers
Method
getGame_rate
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:59
↓ 1 callers
Method
getGame_user
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:55
↓ 1 callers
Method
getInternalTimerService
Returns a {@link InternalTimerService} that can be used to query current processing time and event time and to set timers. An operator can have severa
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:587
↓ 1 callers
Method
getKey
(Tuple2<String, Integer> tuple2)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_01_StreamWordCount.java:67
↓ 1 callers
Method
getLose_user
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:91
↓ 1 callers
Method
getMaxTimestamp
用于计算移除的时间截止点
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:240
↓ 1 callers
Method
getReducingState
(ReducingStateDescriptor<T> stateProperties)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:725
↓ 1 callers
Method
getServer_id
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:51
↓ 1 callers
Method
getWIN_MONEY
()
flink_course/src/main/java/tmp/pojos/UserSlotGame.java:94
↓ 1 callers
Method
getWin_user
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:87
↓ 1 callers
Method
isElementLate
Decide if a record is currently late, based on current watermark and allowed lateness. @param element The element to check @return The element for wh
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:622
↓ 1 callers
Method
join
(Tuple2<String, String> t1, Tuple3<String, String, String> t2)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_15_StreamCoGroup_Join_Demo.java:116
↓ 1 callers
Method
map
(String s)
flink_course/src/main/java/tmp/FlinkKafkaDemo.java:108
↓ 1 callers
Method
map
(String s)
flink_course/src/main/java/cn/doitedu/flink/task/Mapper1.java:5
↓ 1 callers
Method
map
要让flink来帮助管理的状态数据 ,那就不要自己定义一个变量 而是要从flink的api中去获取一个状态管理器,用这个状态管理器来进行数据的增删改查等操作 这种状态: 叫做 托管状态 ! (flink state)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_22_StateBasic_Demo.java:35
↓ 1 callers
Method
map
(String value)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_25_State_DataStructure_Demo.java:124
↓ 1 callers
Method
map
正常的MapFunction的处理逻辑方法 @param value The input value. @return @throws Exception
flink_course/src/main/java/cn/doitedu/flink/java/demos/_23_State_OperatorState_Demo.java:72
↓ 1 callers
Method
map
(String element)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_28_ToleranceSideToSideTest.java:155
↓ 1 callers
Method
map
(String value)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_24_State_KeyedState_Demo.java:68
↓ 1 callers
Method
markActive
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/TimestampsAndWatermarksOperator.java:182
↓ 1 callers
Method
onEvent
(T event, long eventTimestamp, WatermarkOutput output)
flink_course/src/main/java/org/apache/flink/api/common/eventtime/BoundedOutOfOrdernessWatermarks.java:62
↓ 1 callers
Method
onEventTime
(InternalTimer<K, W> timer)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:458
↓ 1 callers
Method
onMerge
(Collection<W> mergedWindows)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:950
↓ 1 callers
Method
onProcessingTime
(InternalTimer<K, W> timer)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:508
↓ 1 callers
Method
open
source组件初始化 @param parameters @throws Exception
flink_course/src/main/java/cn/doitedu/flink/java/demos/_06_CustomSourceFunction.java:106
↓ 1 callers
Method
reduce
(Tuple2<String, Integer> value1, Tuple2<String, Integer> value2)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/Exercise.java:167
↓ 1 callers
Method
setBet_game_num
(long bet_game_num)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:151
↓ 1 callers
Method
setChainingStrategy
(ChainingStrategy strategy)
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:527
↓ 1 callers
Method
setCurrentKey
(Object key)
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:508
↓ 1 callers
Method
setFree_game_sin_money
(long free_game_sin_money)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:135
↓ 1 callers
Method
setGame_id
(long game_id)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:103
↓ 1 callers
Method
setGame_user
(long game_user)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:111
↓ 1 callers
Method
setMain_game_win_money
(long main_game_win_money)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:131
↓ 1 callers
Method
setRandomEvent
(LogBeanWrapper wrapperBean,Collector collector)
datagen/src/main/java/cn/doitedu/module/LogRunnable.java:109
↓ 1 callers
Method
setServer_id
(long server_id)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:107
↓ 1 callers
Method
setTotal_bet_money
(long total_bet_money)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:119
↓ 1 callers
Method
setTotal_game_num
(long total_game_num)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:127
↓ 1 callers
Method
setTotal_win_money
(long total_win_money)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:123
↓ 1 callers
Method
sideOutput
Write skipped late arriving element to SideOutput. @param element skipped late arriving element to side output
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:589
↓ 1 callers
Method
toString
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:1006
↓ 1 callers
Method
userToWrapper
(HashMap<String,LogBean> users)
datagen/src/main/java/cn/doitedu/module/UserUtils.java:90
Method
AbstractPerWindowStateStore
( KeyedStateBackend<?> keyedStateBackend, ExecutionConfig executionConfig)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:697
Method
BoundedOutOfOrdernessWatermarks
Creates a new watermark generator with the given out-of-orderness bound. @param maxOutOfOrderness The bound for the out-of-orderness of the event tim
flink_course/src/main/java/org/apache/flink/api/common/eventtime/BoundedOutOfOrdernessWatermarks.java:50
Method
CollectorKafkaImpl
(String topicName)
datagen/src/main/java/cn/doitedu/module/CollectorKafkaImpl.java:15
Method
ConsumeRunnable
(ConcurrentHashMap<Long, String> guidMap)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习_消费者.java:44
Method
ConsumeRunnableBitmap
(RoaringBitmap bitmap)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习_消费者_Bitmap.java:44
Method
ConsumeRunnableBloomFilter
()
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习_消费者_判重.java:35
Method
Context
(K key, W window)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:832
Method
LogBeanWrapper
(LogBean logBean,String sessionId,long lastTime)
datagen/src/main/java/cn/doitedu/module/LogBeanWrapper.java:22
Method
LogRunnable
@param userPart @param collector @param slowFactor 速度因子,事件间隔时长=RandomUtils.nextInt(50+10 slowFactor,100 slowFactor)
datagen/src/main/java/cn/doitedu/module/LogRunnable.java:22
Method
MergingWindowStateStore
( KeyedStateBackend<?> keyedStateBackend, ExecutionConfig executionConfig)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:708
Method
MyDataGen
()
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习.java:64
Method
MyEventTimeTrigger
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:102
Method
MyTimeEvictor
(long windowSize)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:173
Method
MysqlUser
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:22
Method
PerWindowStateStore
( KeyedStateBackend<?> keyedStateBackend, ExecutionConfig executionConfig)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:751
Method
StatisticBitmapTask
(RoaringBitmap bitmap)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习_消费者_Bitmap.java:84
Method
StatisticTask
(ConcurrentHashMap<Long, String> guidMap)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习_消费者.java:80
Method
Timer
(long timestamp, K key, W window)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:971
Method
TimestampsAndWatermarksOperator
( WatermarkStrategy<T> watermarkStrategy, boolean emitProgressiveWatermarks)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/TimestampsAndWatermarksOperator.java:69
Method
UserSlotGame
()
flink_course/src/main/java/tmp/pojos/UserSlotGame.java:32
Method
WatermarkEmitter
(Output<?> output)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/TimestampsAndWatermarksOperator.java:154
Method
WindowContext
(W window)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:773
Method
WindowOperator
Creates a new {@code WindowOperator} based on the given policies and user functions.
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:179
Method
accept
(PreparedStatement stmt, EventUserInfo eventUserInfo)
flink_course/src/main/java/cn/doitedu/flink/exercise/Exercise_1.java:154
Method
accept
(PreparedStatement preparedStatement, EventLog eventLog)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_11_JdbcSinkOperator_Demo1.java:47
Method
accept
(PreparedStatement preparedStatement, String s)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_28_ToleranceSideToSideTest.java:113
Method
accumulate
进来输入数据后,如何更新累加器 @param accumulator @param score
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo22_CustomAggregateFunction.java:82
Method
canMerge
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:148
Method
cancel
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_06_CustomSourceFunction.java:78
Method
cancel
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_06_CustomSourceFunction.java:91
Method
cancel
job取消调用的方法
flink_course/src/main/java/cn/doitedu/flink/java/demos/_06_CustomSourceFunction.java:151
Method
cancel
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_06_CustomSourceFunction.java:173
Method
clear
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:787
Method
close
()
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:328
Method
close
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:281
Method
close
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_06_CustomSourceFunction.java:184
Method
close
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_17_ProcessFunctions_Demo.java:71
Method
collect
(String logdata)
datagen/src/main/java/cn/doitedu/module/CollectorConsoleImpl.java:4
Method
collect
(String logdata)
datagen/src/main/java/cn/doitedu/module/CollectorKafkaImpl.java:27
Method
configure
(Map<String, ?> configs)
kafka_course/src/main/java/cn/doitedu/kafka/MyPartitioner.java:21
Method
createAccumulator
()
flink_course/src/main/java/tmp/FlinkKafkaDemo.java:139
Method
createAccumulator
初始化累加器 @return
flink_course/src/main/java/cn/doitedu/flink/java/demos/_20_Window_Api_Demo1.java:90
Method
createAccumulator
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_25_State_DataStructure_Demo.java:96
Method
createAccumulator
创建累加器 @return
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo22_CustomAggregateFunction.java:66
Method
createAccumulator
()
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo24_TableAggregateFunction.java:96
Method
createAccumulator
()
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo24_TableAggregateFunction2.java:99
Method
createAccumulator
()
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/Exercise.java:179
← previous
next →
101–200 of 478, ranked by callers