MCPcopy Create free account

hub / github.com/PacktPublishing/Building-Big-Data-Pipelines-with-Apache-Beam / functions

Functions408 in github.com/PacktPublishing/Building-Big-Data-Pipelines-with-Apache-Beam

↓ 1 callersFunctionmove
(pos, direction, speed, t)
chapter4/bin/emit_tracks.py:16
↓ 1 callersMethodnewContext
( WindowFn<T, BoundedWindow> windowFn, Instant timestamp)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter2.java:421
↓ 1 callersMethodnewTimestampPolicy
()
chapter5/src/main/java/com/packtpub/beam/chapter5/SchemaSportTracker.java:192
↓ 1 callersMethodnewTimestampPolicy
()
chapter2/src/main/java/com/packtpub/beam/chapter2/SportTracker.java:172
↓ 1 callersMethodnewTimestampPolicy
()
util/src/main/java/com/packtpub/beam/util/ReadPositionsFromKafka.java:76
↓ 1 callersMethodnewValue
(Instant now, int pos)
chapter3/src/test/java/com/packtpub/beam/chapter3/DroppableDataFilter2Test.java:138
↓ 1 callersMethodofProcessingTime
(Duration duration)
chapter1/src/main/java/com/packtpub/beam/chapter1/WithReadDelay.java:33
↓ 1 callersMethodoutputLineAsPositionWithCreateTimestamp
( String line, KafkaProducer<String, String> producer, String outputTopic)
util/src/main/java/com/packtpub/beam/util/WritePositionsToKafka.java:62
↓ 1 callersMethodoutputWindowsWithinInterval
( OutputReceiver<BoundedWindow> output, Instant minStamp, Instant maxStamp)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter2.java:407
↓ 1 callersMethodparseArgs
(String[] args)
chapter5/src/main/java/com/packtpub/beam/chapter5/SchemaSportTracker.java:203
↓ 1 callersMethodparseArgs
(String[] args)
chapter5/src/main/java/com/packtpub/beam/chapter5/SQLMaxWordLength.java:108
↓ 1 callersMethodparseArgs
(String[] args)
chapter5/src/main/java/com/packtpub/beam/chapter5/SQLSportTrackerMotivation.java:130
↓ 1 callersMethodparseArgs
(String[] args)
chapter2/src/main/java/com/packtpub/beam/chapter2/MaxWordLength.java:85
↓ 1 callersMethodparseArgs
(String[] args)
chapter2/src/main/java/com/packtpub/beam/chapter2/AverageWordLength.java:95
↓ 1 callersMethodparseArgs
(String[] args)
chapter2/src/main/java/com/packtpub/beam/chapter2/SlidingWindowWordLength.java:72
↓ 1 callersMethodparseArgs
(String[] args)
chapter2/src/main/java/com/packtpub/beam/chapter2/SportTracker.java:183
↓ 1 callersMethodparseArgs
(String[] args)
chapter4/src/main/java/com/packtpub/beam/chapter4/SportTrackerMotivationUsingJoinLibrary.java:96
↓ 1 callersMethodparseArgs
(String[] args)
chapter4/src/main/java/com/packtpub/beam/chapter4/SportTrackerMotivationUsingSideInputs.java:102
↓ 1 callersMethodparseArgs
(String[] args)
chapter4/src/main/java/com/packtpub/beam/chapter4/SportTrackerMotivationUsingCoGBK.java:101
↓ 1 callersMethodparseArgs
(String[] args)
chapter4/src/main/java/com/packtpub/beam/chapter4/SportTrackerMotivationUsingOwnJoin.java:113
↓ 1 callersMethodparseArgs
(String[] args)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDo.java:139
↓ 1 callersMethodparseArgs
(String[] args)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter2.java:227
↓ 1 callersMethodparseArgs
(String[] args)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoBatch.java:184
↓ 1 callersMethodparseArgs
(String[] args)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoStateful.java:328
↓ 1 callersMethodparseArgs
(String[] args)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter.java:172
↓ 1 callersMethodprocessNewFiles
( RestrictionTracker<DirectoryWatchRestriction, String> tracker, ManualWatermarkEstimator<Inst
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:231
↓ 1 callersMethodreadInput
(Pipeline pipeline, Params params)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDo.java:72
↓ 1 callersMethodreadInput
(Pipeline pipeline, Params params)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter2.java:94
↓ 1 callersMethodreadInput
(Pipeline pipeline, Params params)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoBatch.java:81
↓ 1 callersMethodreadInput
(Pipeline pipeline, Params params)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoStateful.java:114
↓ 1 callersMethodreadInput
(Pipeline pipeline, Params params)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter.java:85
↓ 1 callersFunctionreadToStream
(filename)
chapter6/src/main/python/first_streaming_pipeline.py:15
↓ 1 callersMethodresolve
(Request request, StreamObserver<Response> responseObserver)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCService.java:27
↓ 1 callersMethodrewindow
( PCollection<T> input, WindowingStrategy<T, BoundedWindow> strategy)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter2.java:206
↓ 1 callersMethodrewindow
( PCollection<String> input, WindowingStrategy<String, BoundedWindow> strategy)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter.java:151
↓ 1 callersMethodrunRpc
(int port)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoBatch.java:106
↓ 1 callersMethodrunRpc
(int port)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoStateful.java:139
↓ 1 callersMethodseekFileToStartLine
(RandomAccessFile file, long position)
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:333
↓ 1 callersMethodsetNewLoopTimer
(Timer timer, Instant currentStamp, long minTsToKeep)
util/src/main/java/com/packtpub/beam/util/ToMetric.java:165
↓ 1 callersMethodstoreResult
(PCollection<KV<String, Integer>> result, Params params)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDo.java:83
↓ 1 callersMethodstoreResult
(PCollection<KV<String, Integer>> result, Params params)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoBatch.java:92
↓ 1 callersMethodstoreResult
(PCollection<KV<String, Integer>> result, Params params)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoStateful.java:125
↓ 1 callersMethodtoDelete
()
chapter4/src/main/java/com/packtpub/beam/chapter4/JoinRow.java:70
↓ 1 callersFunctionusage
()
chapter6/src/main/python/first_pipeline.py:8
↓ 1 callersFunctionusage
()
chapter6/src/main/python/max_word_length.py:14
↓ 1 callersFunctionusage
()
chapter6/src/main/python/rpc_par_do.py:19
↓ 1 callersFunctionusage
()
chapter6/src/main/python/sport_tracker.py:13
↓ 1 callersFunctionusage
()
chapter6/src/main/python/sport_tracker_motivation.py:14
↓ 1 callersFunctionusage
()
chapter6/src/main/python/first_streaming_pipeline.py:11
↓ 1 callersFunctionusage
()
chapter4/bin/emit_tracks.py:12
↓ 1 callersMethodusage
()
chapter7/src/main/java/com/packtpub/beam/chapter7/PiSampler.java:48
↓ 1 callersMethodusage
()
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:74
↓ 1 callersMethodusage
()
util/src/main/java/com/packtpub/beam/util/WritePositionsToKafka.java:29
↓ 1 callersMethodwithRandomFactory
(RandomFactory randomFactory)
chapter7/src/main/java/com/packtpub/beam/chapter7/PiSampler.java:88
↓ 1 callersMethodwithWatermarkFn
(SerializableFunction<KV<String, Instant>, Instant> watermarkFn)
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:105
↓ 1 callersMethodwriteStdinToKafkaWithTimestamps
(String bootstrapServer, String outputTopic)
util/src/main/java/com/packtpub/beam/util/WritePositionsToKafka.java:45
MethodAutoSchemaPosition
(double latitude, double longitude)
chapter5/src/test/java/com/packtpub/beam/chapter5/TestSchemaInference.java:59
MethodBatchRpcDoFn
(String hostname, int port)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoBatch.java:130
MethodBatchRpcDoFnStateful
(int maxBatchSize, Duration maxBatchWait, String hostname, int port)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoStateful.java:212
MethodCacheAndUpdateJoinKeyFn
(Coder<JoinRow<K, J, V>> rowCoder)
chapter4/src/main/java/com/packtpub/beam/chapter4/StreamingInnerJoin.java:164
MethodChapter1DemoTest
()
chapter1/src/test/java/com/packtpub/beam/chapter1/Chapter1DemoTest.java:47
MethodDebugOutput
(String name)
util/src/main/java/com/packtpub/beam/util/DebugOutput.java:33
MethodDelayFn
(Duration delay)
chapter1/src/main/java/com/packtpub/beam/chapter1/WithReadDelay.java:58
MethodDirectoryWatchFn
(SerializableFunction<KV<String, Instant>, Instant> watermarkFn)
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:209
MethodDirectoryWatchRestrictionTracker
(DirectoryWatchRestriction restriction)
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:173
MethodJoinResultCoder
( Coder<K> keyCoder, Coder<J> joinKeyCoder, Coder<L> leftValueCoder, Coder<R>
chapter4/src/main/java/com/packtpub/beam/chapter4/JoinResult.java:50
MethodJoinRowCoder
(Coder<K> keyCoder, Coder<J> joinKeyCoder, Coder<V> valueCoder)
chapter4/src/main/java/com/packtpub/beam/chapter4/JoinRow.java:44
MethodLoopTimerForWindowLabels
( Duration loopDuration, Instant startingTime, WindowingStrategy<KV<String, String>, ?
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter2.java:368
MethodMapToLines
(boolean includeKey)
util/src/main/java/com/packtpub/beam/util/MapToLines.java:37
MethodMetric
(double length, long duration)
util/src/main/java/com/packtpub/beam/util/ToMetric.java:63
MethodMetricAccumulator
()
chapter4/src/main/java/com/packtpub/beam/chapter4/ComputeAverage.java:83
MethodPiSampler
(long numSamples, long parallelism)
chapter7/src/main/java/com/packtpub/beam/chapter7/PiSampler.java:76
MethodReadPositionsFromKafka
(String bootstrapServer, String inputTopic)
util/src/main/java/com/packtpub/beam/util/ReadPositionsFromKafka.java:42
MethodRpcDoFn
(String hostname, int port)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDo.java:112
MethodSplitDroppableDataFn
( Duration allowedLateness, Coder<BoundedWindow> windowCoder, TupleTag<String> mainOut
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter2.java:264
MethodSplitDroppableDataFn
( Duration allowedLateness, TupleTag<String> mainOutput, TupleTag<String> droppableOutput)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter.java:201
MethodStreamInnerJoinDoFn
(Coder<K> keyCoder, Coder<L> leftValueCoder, Coder<R> rightValueCoder)
chapter4/src/main/java/com/packtpub/beam/chapter4/StreamingInnerJoin.java:70
MethodStreamingFileRead
(String directoryPath)
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:350
MethodStreamingInnerJoin
( TupleTag<L> leftHandTag, TupleTag<R> rightHandTag, SerializableFunction<L, K> leftKeyExtra
chapter4/src/main/java/com/packtpub/beam/chapter4/StreamingInnerJoin.java:198
MethodUtils
()
util/src/main/java/com/packtpub/beam/util/Utils.java:51
MethodValueWithTimestampCoder
(Coder<T> valueCoder)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoStateful.java:161
MethodWithReadDelay
(Duration delay)
chapter1/src/main/java/com/packtpub/beam/chapter1/WithReadDelay.java:39
MethodWithStringSchema
(String name)
util/src/main/java/com/packtpub/beam/util/WithStringSchema.java:34
MethodWriteNotificationsToKafka
(String bootstrapServer, String outputTopic)
util/src/main/java/com/packtpub/beam/util/WriteNotificationsToKafka.java:32
Method__enter__
(self)
chapter6/src/main/python/rpc_par_do.py:34
Method__exit__
(self, exc_type, exc_val, exc_tb)
chapter6/src/main/python/rpc_par_do.py:39
Method__init__
(self, port=None)
chapter6/src/main/python/rpc_par_do.py:31
Method__init__
(self, address, batchSize, maxWaitTime)
chapter6/src/main/python/rpc_par_do.py:71
Method__init__
(self, duration)
chapter6/src/main/python/sport_tracker_motivation.py:54
MethodaddInput
(List<Row> mutableAccumulator, Row input)
chapter5/src/main/java/com/packtpub/beam/chapter5/SchemaSportTracker.java:231
MethodaddInput
(AverageAccumulator accumulator, String input)
chapter2/src/main/java/com/packtpub/beam/chapter2/AverageWordLength.java:129
MethodaddInput
(MetricAccumulator accumulator, Metric input)
chapter4/src/main/java/com/packtpub/beam/chapter4/ComputeAverage.java:55
Methodadd_input
(self, acc, element)
chapter6/src/main/python/sport_tracker_motivation.py:25
Methodapply
(String input)
chapter5/src/main/java/com/packtpub/beam/chapter5/SQLMaxWordLength.java:102
FunctionasMotivation
(x)
chapter6/src/main/python/sport_tracker_motivation.py:134
MethodcheckDone
()
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:196
MethodcomputeRawMetrics
( KV<String, Iterable<Position>> stringIterableKV)
chapter2/src/main/java/com/packtpub/beam/chapter2/SportTracker.java:148
MethodcreateAccumulator
()
chapter5/src/main/java/com/packtpub/beam/chapter5/SchemaSportTracker.java:226
MethodcreateAccumulator
()
chapter2/src/main/java/com/packtpub/beam/chapter2/AverageWordLength.java:124
MethodcreateAccumulator
()
chapter4/src/main/java/com/packtpub/beam/chapter4/ComputeAverage.java:44
← previousnext →101–200 of 408, ranked by callers