MCPcopy Create free account

hub / github.com/coderblack/doit30_flink / functions

Functions478 in github.com/coderblack/doit30_flink

↓ 1 callersMethoddeleteCleanupTimer
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 callersMethodextractTimestamp
(UserSlotGame ele, long l)
flink_course/src/main/java/tmp/FlinkKafkaDemo.java:89
↓ 1 callersMethodflatMap
(String s, Collector<Tuple2<String, Integer>> collector)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_01_StreamWordCount.java:49
↓ 1 callersMethodflatMap
(String value, Collector<Tuple2<String, Integer>> out)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_03_StreamBatchWordCount.java:34
↓ 1 callersMethodflatMap
(UserInfo value, Collector<UserFriendInfo> out)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_07_Transformation_Demos.java:72
↓ 1 callersMethodflatMap
(String value, Collector<String> out)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_18_ChannalSelector_Partitioner_Demo.java:27
↓ 1 callersMethodflatMap
(String value, Collector<String> out)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/KeyedStateDemo.java:35
↓ 1 callersMethodgenBatchToConsole
(List<LogBeanWrapper> wrapperedUsers, int threads , Collector collector)
datagen/src/main/java/cn/doitedu/ActionLogAutoGen.java:81
↓ 1 callersMethodgenData
()
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习.java:76
↓ 1 callersMethodgetAggregatingState
( AggregatingStateDescriptor<IN, ACC, OUT> stateProperties)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:731
↓ 1 callersMethodgetContainingTask
()
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:404
↓ 1 callersMethodgetGame_id
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:47
↓ 1 callersMethodgetGame_rate
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:59
↓ 1 callersMethodgetGame_user
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:55
↓ 1 callersMethodgetInternalTimerService
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 callersMethodgetKey
(Tuple2<String, Integer> tuple2)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_01_StreamWordCount.java:67
↓ 1 callersMethodgetLose_user
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:91
↓ 1 callersMethodgetMaxTimestamp
用于计算移除的时间截止点
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:240
↓ 1 callersMethodgetReducingState
(ReducingStateDescriptor<T> stateProperties)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:725
↓ 1 callersMethodgetServer_id
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:51
↓ 1 callersMethodgetWIN_MONEY
()
flink_course/src/main/java/tmp/pojos/UserSlotGame.java:94
↓ 1 callersMethodgetWin_user
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:87
↓ 1 callersMethodisElementLate
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 callersMethodjoin
(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 callersMethodmap
(String s)
flink_course/src/main/java/tmp/FlinkKafkaDemo.java:108
↓ 1 callersMethodmap
(String s)
flink_course/src/main/java/cn/doitedu/flink/task/Mapper1.java:5
↓ 1 callersMethodmap
要让flink来帮助管理的状态数据 ,那就不要自己定义一个变量 而是要从flink的api中去获取一个状态管理器,用这个状态管理器来进行数据的增删改查等操作 这种状态: 叫做 托管状态 ! (flink state)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_22_StateBasic_Demo.java:35
↓ 1 callersMethodmap
(String value)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_25_State_DataStructure_Demo.java:124
↓ 1 callersMethodmap
正常的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 callersMethodmap
(String element)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_28_ToleranceSideToSideTest.java:155
↓ 1 callersMethodmap
(String value)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_24_State_KeyedState_Demo.java:68
↓ 1 callersMethodmarkActive
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/TimestampsAndWatermarksOperator.java:182
↓ 1 callersMethodonEvent
(T event, long eventTimestamp, WatermarkOutput output)
flink_course/src/main/java/org/apache/flink/api/common/eventtime/BoundedOutOfOrdernessWatermarks.java:62
↓ 1 callersMethodonEventTime
(InternalTimer<K, W> timer)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:458
↓ 1 callersMethodonMerge
(Collection<W> mergedWindows)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:950
↓ 1 callersMethodonProcessingTime
(InternalTimer<K, W> timer)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:508
↓ 1 callersMethodopen
source组件初始化 @param parameters @throws Exception
flink_course/src/main/java/cn/doitedu/flink/java/demos/_06_CustomSourceFunction.java:106
↓ 1 callersMethodreduce
(Tuple2<String, Integer> value1, Tuple2<String, Integer> value2)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/Exercise.java:167
↓ 1 callersMethodsetBet_game_num
(long bet_game_num)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:151
↓ 1 callersMethodsetChainingStrategy
(ChainingStrategy strategy)
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:527
↓ 1 callersMethodsetCurrentKey
(Object key)
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:508
↓ 1 callersMethodsetFree_game_sin_money
(long free_game_sin_money)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:135
↓ 1 callersMethodsetGame_id
(long game_id)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:103
↓ 1 callersMethodsetGame_user
(long game_user)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:111
↓ 1 callersMethodsetMain_game_win_money
(long main_game_win_money)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:131
↓ 1 callersMethodsetRandomEvent
(LogBeanWrapper wrapperBean,Collector collector)
datagen/src/main/java/cn/doitedu/module/LogRunnable.java:109
↓ 1 callersMethodsetServer_id
(long server_id)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:107
↓ 1 callersMethodsetTotal_bet_money
(long total_bet_money)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:119
↓ 1 callersMethodsetTotal_game_num
(long total_game_num)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:127
↓ 1 callersMethodsetTotal_win_money
(long total_win_money)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:123
↓ 1 callersMethodsideOutput
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 callersMethodtoString
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:1006
↓ 1 callersMethoduserToWrapper
(HashMap<String,LogBean> users)
datagen/src/main/java/cn/doitedu/module/UserUtils.java:90
MethodAbstractPerWindowStateStore
( KeyedStateBackend<?> keyedStateBackend, ExecutionConfig executionConfig)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:697
MethodBoundedOutOfOrdernessWatermarks
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
MethodCollectorKafkaImpl
(String topicName)
datagen/src/main/java/cn/doitedu/module/CollectorKafkaImpl.java:15
MethodConsumeRunnable
(ConcurrentHashMap<Long, String> guidMap)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习_消费者.java:44
MethodConsumeRunnableBitmap
(RoaringBitmap bitmap)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习_消费者_Bitmap.java:44
MethodConsumeRunnableBloomFilter
()
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习_消费者_判重.java:35
MethodContext
(K key, W window)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:832
MethodLogBeanWrapper
(LogBean logBean,String sessionId,long lastTime)
datagen/src/main/java/cn/doitedu/module/LogBeanWrapper.java:22
MethodLogRunnable
@param userPart @param collector @param slowFactor 速度因子,事件间隔时长=RandomUtils.nextInt(50+10 slowFactor,100 slowFactor)
datagen/src/main/java/cn/doitedu/module/LogRunnable.java:22
MethodMergingWindowStateStore
( KeyedStateBackend<?> keyedStateBackend, ExecutionConfig executionConfig)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:708
MethodMyDataGen
()
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习.java:64
MethodMyEventTimeTrigger
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:102
MethodMyTimeEvictor
(long windowSize)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:173
MethodMysqlUser
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:22
MethodPerWindowStateStore
( KeyedStateBackend<?> keyedStateBackend, ExecutionConfig executionConfig)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:751
MethodStatisticBitmapTask
(RoaringBitmap bitmap)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习_消费者_Bitmap.java:84
MethodStatisticTask
(ConcurrentHashMap<Long, String> guidMap)
kafka_course/src/main/java/cn/doitedu/kafka/Kafka编程练习_消费者.java:80
MethodTimer
(long timestamp, K key, W window)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:971
MethodTimestampsAndWatermarksOperator
( WatermarkStrategy<T> watermarkStrategy, boolean emitProgressiveWatermarks)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/TimestampsAndWatermarksOperator.java:69
MethodUserSlotGame
()
flink_course/src/main/java/tmp/pojos/UserSlotGame.java:32
MethodWatermarkEmitter
(Output<?> output)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/TimestampsAndWatermarksOperator.java:154
MethodWindowContext
(W window)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:773
MethodWindowOperator
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
Methodaccept
(PreparedStatement stmt, EventUserInfo eventUserInfo)
flink_course/src/main/java/cn/doitedu/flink/exercise/Exercise_1.java:154
Methodaccept
(PreparedStatement preparedStatement, EventLog eventLog)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_11_JdbcSinkOperator_Demo1.java:47
Methodaccept
(PreparedStatement preparedStatement, String s)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_28_ToleranceSideToSideTest.java:113
Methodaccumulate
进来输入数据后,如何更新累加器 @param accumulator @param score
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo22_CustomAggregateFunction.java:82
MethodcanMerge
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:148
Methodcancel
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_06_CustomSourceFunction.java:78
Methodcancel
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_06_CustomSourceFunction.java:91
Methodcancel
job取消调用的方法
flink_course/src/main/java/cn/doitedu/flink/java/demos/_06_CustomSourceFunction.java:151
Methodcancel
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_06_CustomSourceFunction.java:173
Methodclear
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:787
Methodclose
()
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:328
Methodclose
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:281
Methodclose
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_06_CustomSourceFunction.java:184
Methodclose
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_17_ProcessFunctions_Demo.java:71
Methodcollect
(String logdata)
datagen/src/main/java/cn/doitedu/module/CollectorConsoleImpl.java:4
Methodcollect
(String logdata)
datagen/src/main/java/cn/doitedu/module/CollectorKafkaImpl.java:27
Methodconfigure
(Map<String, ?> configs)
kafka_course/src/main/java/cn/doitedu/kafka/MyPartitioner.java:21
MethodcreateAccumulator
()
flink_course/src/main/java/tmp/FlinkKafkaDemo.java:139
MethodcreateAccumulator
初始化累加器 @return
flink_course/src/main/java/cn/doitedu/flink/java/demos/_20_Window_Api_Demo1.java:90
MethodcreateAccumulator
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_25_State_DataStructure_Demo.java:96
MethodcreateAccumulator
创建累加器 @return
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo22_CustomAggregateFunction.java:66
MethodcreateAccumulator
()
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo24_TableAggregateFunction.java:96
MethodcreateAccumulator
()
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo24_TableAggregateFunction2.java:99
MethodcreateAccumulator
()
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/Exercise.java:179
← previousnext →101–200 of 478, ranked by callers