MCPcopy Create free account

hub / github.com/coderblack/doit30_flink / functions

Functions478 in github.com/coderblack/doit30_flink

↓ 115 callersMethodof
(Time windowSize)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:251
↓ 57 callersMethodcollect
(String logdata)
datagen/src/main/java/cn/doitedu/module/Collector.java:4
↓ 37 callersMethodcreate
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:161
↓ 35 callersMethodget
()
flink_course/src/main/java/cn/doitedu/flink/java/demos/_11_JdbcSinkOperator_Demo1.java:105
↓ 32 callersMethodmap
(String value)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_26_State_TTL_Demo.java:132
↓ 23 callersMethodmap
(String value)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/Exercise.java:78
↓ 15 callersMethodgetRuntimeContext
Returns a context that allows the operator to query information about the execution and also to interact with systems such as broadcast variables and
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:432
↓ 12 callersMethodequals
(Object o)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:982
↓ 11 callersMethodgetValue
()
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo25_MetricDemos.java:60
↓ 11 callersMethodprocess
(String key, ProcessWindowFunction<EventBean2, Tuple3<String,String,Double>, String, TimeWindow>.Context conte
flink_course/src/main/java/cn/doitedu/flink/java/demos/_20_Window_Api_Demo1.java:234
↓ 10 callersMethodadd
(UserSlotGame value, MysqlUser accumulator)
flink_course/src/main/java/tmp/FlinkKafkaDemo.java:144
↓ 9 callersMethodclear
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:955
↓ 9 callersMethodcurrentWatermark
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:796
↓ 8 callersMethodclose
()
kafka_course/src/main/java/cn/doitedu/kafka/MyPartitioner.java:16
↓ 7 callersMethodprocess
(String key, ProcessWindowFunction<EventBean, String, String, TimeWindow>.Context context, Iterable<EventBean>
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/Exercise.java:211
↓ 6 callersMethodadd
滚动聚合的逻辑(拿到一条数据,如何去更新累加器) @param value The value to add @param accumulator The accumulator to add the value to @return
flink_course/src/main/java/cn/doitedu/flink/java/demos/_20_Window_Api_Demo1.java:101
↓ 6 callersMethodaggregate
TODO 累加等聚合逻辑 @param value @param accumulator GAME_ID,SERVER_ID,USER_ID max(LOG_DATE) as LOG_DATE sum(BET_MONEY) as total_bet_money sum(WIN_MONEY) as
flink_course/src/main/java/tmp/FlinkKafkaDemo.java:261
↓ 6 callersMethodfilter
(EventBean value)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/Exercise.java:140
↓ 6 callersMethodgetLog_date
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:43
↓ 6 callersMethodgetProcessingTimeService
Returns the {@link ProcessingTimeService} responsible for getting the current processing time and registering timers.
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:451
↓ 6 callersMethodgetState
(ValueStateDescriptor<T> stateProperties)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:713
↓ 5 callersMethodadd
(Integer value, Tuple2<Integer, Integer> accumulator)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_25_State_DataStructure_Demo.java:101
↓ 5 callersMethodcurrentProcessingTime
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:791
↓ 5 callersMethodgetBUFF0
()
flink_course/src/main/java/tmp/pojos/UserSlotGame.java:122
↓ 5 callersMethodgetExecutionConfig
Gets the execution config defined on the execution environment of the job to which this operator belongs. @return The job's execution config.
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:396
↓ 5 callersMethodgetKeyedStateBackend
()
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:437
↓ 5 callersMethodgetListState
(ListStateDescriptor<T> stateProperties)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:719
↓ 5 callersMethodpartition
(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster)
kafka_course/src/main/java/cn/doitedu/kafka/MyPartitioner.java:9
↓ 4 callersMethodcleanupTime
Returns the cleanup time for a window, which is {@code window.maxTimestamp + allowedLateness}. In case this leads to a value greater than {@link Long#
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:672
↓ 4 callersMethodemitWindowContents
Emits the contents of the given window using the {@link InternalWindowFunction}.
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:576
↓ 4 callersMethodgetBonus_game_sin_money
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:83
↓ 4 callersMethodgetCurrentKey
()
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:512
↓ 4 callersMethodgetKey
(EventBean value)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/Exercise.java:98
↓ 4 callersMethodgetOperatorID
()
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:652
↓ 4 callersMethodmap
(String s)
flink_course/src/main/java/cn/doitedu/flink/task/Mapper2.java:4
↓ 4 callersMethodoutput
(OutputTag<X> outputTag, X value)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:812
↓ 4 callersMethodprocess
(String s, ProcessWindowFunction<Tuple2<String, Long>, String, String, TimeWindow>.Context context, Iterable<T
flink_course/src/main/java/cn/doitedu/flink/TestWindow.java:40
↓ 3 callersMethodapply
@param key 本次传给咱们的窗口是属于哪个key的 @param window 本次传给咱们的窗口的各种元信息(比如本窗口的起始时间,结束时间) @param input 本次传给咱们的窗口中所有数据的迭代器 @param out 结果数据输出器 @throws Exception
flink_course/src/main/java/cn/doitedu/flink/java/demos/_20_Window_Api_Demo1.java:295
↓ 3 callersMethodemitWatermark
(Watermark watermark)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/TimestampsAndWatermarksOperator.java:159
↓ 3 callersMethodgetGAME_TYPE
()
flink_course/src/main/java/tmp/pojos/UserSlotGame.java:82
↓ 3 callersMethodgetLOG_DATE
()
flink_course/src/main/java/tmp/pojos/UserSlotGame.java:146
↓ 3 callersMethodgetMergingWindowSet
Retrieves the {@link MergingWindowSet} for the currently active key. The caller must ensure that the correct key is set in the state backend. <p>The
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:600
↓ 3 callersMethodgetMetricGroup
()
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:251
↓ 3 callersMethodgetOrCreateKeyedState
( TypeSerializer<N> namespaceSerializer, StateDescriptor<S, T> stateDescriptor)
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:468
↓ 3 callersMethodgetUserCodeClassloader
()
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:408
↓ 3 callersMethodprocessWatermark
(Watermark mark)
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:601
↓ 3 callersMethodregisterEventTimeTimer
(long time)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:923
↓ 3 callersMethodreportOrForwardLatencyMarker
(LatencyMarker marker)
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:555
↓ 3 callersMethodtoString
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:155
↓ 2 callersMethodaccumulate
累加更新逻辑 @param acc @param value
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo24_TableAggregateFunction.java:113
↓ 2 callersMethodaccumulate
累加更新逻辑 @param acc @param value
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo24_TableAggregateFunction2.java:116
↓ 2 callersMethodadd
(int i)
flink_course/src/main/java/cn/doitedu/flinksql/demos/Demo25_MetricDemos.java:56
↓ 2 callersMethodapply
(GlobalWindow window, Iterable<EventBean2> values, Collector<String> out)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo2.java:49
↓ 2 callersMethodclear
(TimeWindow window, TriggerContext ctx)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:143
↓ 2 callersMethodclearAllState
Drops all state for the given window and calls {@link Trigger#clear(Window, Trigger.TriggerContext)}. <p>The caller must ensure that the correct key
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:562
↓ 2 callersMethodclose
()
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/Exercise.java:131
↓ 2 callersMethodcompareTo
(Timer<K, W> o)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:977
↓ 2 callersMethoddeleteEventTimeTimer
(long time)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:933
↓ 2 callersMethoddeleteProcessingTimeTimer
(long time)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:928
↓ 2 callersMethodevict
元素移除的核心逻辑
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:208
↓ 2 callersMethodflatMap
(String value, Collector<Tuple2<String, Integer>> out)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_02_BatchWordCount.java:37
↓ 2 callersMethodgetBET_MONEY
()
flink_course/src/main/java/tmp/pojos/UserSlotGame.java:90
↓ 2 callersMethodgetBet_game_num
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:95
↓ 2 callersMethodgetCurrentProcessingTime
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:232
↓ 2 callersMethodgetCurrentWatermark
()
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:842
↓ 2 callersMethodgetFree_game_sin_money
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:79
↓ 2 callersMethodgetGAME_ID
()
flink_course/src/main/java/tmp/pojos/UserSlotGame.java:74
↓ 2 callersMethodgetKeyedStateStore
()
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:516
↓ 2 callersMethodgetMain_game_win_money
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:75
↓ 2 callersMethodgetMapState
(MapStateDescriptor<UK, UV> stateProperties)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:738
↓ 2 callersMethodgetOperatorName
Return the operator name. If the runtime context has been set, then the task name with subtask index is returned. Otherwise, the simple class name is
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:419
↓ 2 callersMethodgetPartitionedState
(StateDescriptor<S, ?> stateDescriptor)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:879
↓ 2 callersMethodgetSERVER_ID
()
flink_course/src/main/java/tmp/pojos/UserSlotGame.java:78
↓ 2 callersMethodgetTotal_bet_money
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:63
↓ 2 callersMethodgetTotal_game_num
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:71
↓ 2 callersMethodgetTotal_win_money
()
flink_course/src/main/java/tmp/pojos/MysqlUser.java:67
↓ 2 callersMethodgetUSER_ID
()
flink_course/src/main/java/tmp/pojos/UserSlotGame.java:66
↓ 2 callersMethodhasTimestamp
(Iterable<TimestampedValue<Object>> elements)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:229
↓ 2 callersMethodisCleanupTime
Returns {@code true} if the given time is the cleanup time for the given window.
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:682
↓ 2 callersMethodisUsingCustomRawKeyedState
Indicates whether or not implementations of this class is writing to the raw keyed state streams on snapshots, using {@link #snapshotState(StateSnapsh
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:309
↓ 2 callersMethodisWindowLate
Returns {@code true} if the watermark is after the end timestamp plus the allowed lateness of the given window.
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:611
↓ 2 callersMethodmap
(String value)
flink_course/src/test/java/cn/doitedu/flink/TestChangelog.java:32
↓ 2 callersMethodonElement
(StreamRecord<IN> element)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:938
↓ 2 callersMethodonPeriodicEmit
(WatermarkOutput output)
flink_course/src/main/java/org/apache/flink/api/common/eventtime/BoundedOutOfOrdernessWatermarks.java:67
↓ 2 callersMethodprocess
(Tuple3<Long, Long, Long> aLong, Context context, Iterable<UserSlotGame> iterable, Collector<MysqlUser> collec
flink_course/src/main/java/tmp/FlinkKafkaDemo.java:182
↓ 2 callersMethodprocessWatermarkStatus
(WatermarkStatus watermarkStatus)
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:629
↓ 2 callersMethodreduce
@param value1 是此前的聚合结果 @param value2 是本次的新数据 @return 更新后的聚合结果 @throws Exception
flink_course/src/main/java/cn/doitedu/flink/java/demos/_07_Transformation_Demos.java:142
↓ 2 callersMethodregisterCleanupTimer
Registers a timer to cleanup the content of the window. @param window the window whose state to discard
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:633
↓ 2 callersMethodregisterProcessingTimeTimer
(long time)
flink_course/src/main/java/org/apache/flink/streaming/runtime/operators/windowing/WindowOperator.java:918
↓ 2 callersMethodsaveUsers
(HashMap<String,LogBean> users)
datagen/src/main/java/cn/doitedu/module/UserUtils.java:75
↓ 2 callersMethodsetBonus_game_sin_money
(long bonus_game_sin_money)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:139
↓ 2 callersMethodsetKeyContextElement
(StreamRecord<T> record, KeySelector<T, ?> selector)
flink_course/src/main/java/org/apache/flink/streaming/api/operators/AbstractStreamOperator.java:500
↓ 2 callersMethodsetLog_date
(String log_date)
flink_course/src/main/java/tmp/pojos/MysqlUser.java:99
↓ 1 callersMethodadd
(EventBean value, Tuple2<String, Integer> accumulator)
flink_course/src/main/java/cn/doitedu/flinksql/fuxi/Exercise.java:184
↓ 1 callersMethodaddNewUsers
(HashMap<String,LogBean> users,int cnt,boolean save)
datagen/src/main/java/cn/doitedu/module/UserUtils.java:36
↓ 1 callersMethodapply
(Long aLong, TimeWindow window, Iterable<Tuple2<EventBean2, Integer>> input, Collector<String> out)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo4.java:74
↓ 1 callersMethodapply
(Long aLong, TimeWindow window, Iterable<Tuple2<EventBean2, Integer>> input, Collector<String> out)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_21_Window_Api_Demo3.java:54
↓ 1 callersMethodclose
组件关闭调用的方法 @throws Exception
flink_course/src/main/java/cn/doitedu/flink/java/demos/_06_CustomSourceFunction.java:160
↓ 1 callersMethodcoGroup
@param first 是协同组中的第一个流的数据 @param second 是协同组中的第二个流的数据 @param out 是处理结果的输出器 @throws Exception
flink_course/src/main/java/cn/doitedu/flink/java/demos/_15_StreamCoGroup_Join_Demo.java:73
↓ 1 callersMethodcompare
(Tuple2<String, Double> tp1, Tuple2<String, Double> tp2)
flink_course/src/main/java/cn/doitedu/flink/java/demos/_20_Window_Api_Demo1.java:258
next →1–100 of 478, ranked by callers