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

↓ 367 callersMethodapply
()
chapter7/src/main/java/com/packtpub/beam/chapter7/PiSampler.java:45
↓ 226 callersMethodof
( J joinKey, @Nullable K leftKey, @Nullable L leftValue, @Nullable K rightKey, @
chapter4/src/main/java/com/packtpub/beam/chapter4/JoinResult.java:31
↓ 165 callersMethodof
(long numSamples, long parallelism)
chapter7/src/main/java/com/packtpub/beam/chapter7/PiSampler.java:54
↓ 123 callersMethodcreate
(String name)
util/src/main/java/com/packtpub/beam/util/DebugOutput.java:27
↓ 62 callersMethodof
(Coder<T> valueCoder)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoStateful.java:154
↓ 49 callersMethodof
()
util/src/main/java/com/packtpub/beam/util/Tokenize.java:25
↓ 34 callersMethodmove
( double latitudeDirection, double longitudeDirection, double speedMeterPerSec, long t
util/src/main/java/com/packtpub/beam/util/Position.java:50
↓ 28 callersMethodrandom
(long stamp)
util/src/main/java/com/packtpub/beam/util/Position.java:42
↓ 27 callersMethoddistance
(Position other)
util/src/main/java/com/packtpub/beam/util/Position.java:46
↓ 19 callersMethodcollect
( J joinKey, K primaryKey, Object primaryValue, K otherKey,
chapter4/src/main/java/com/packtpub/beam/chapter4/StreamingInnerJoin.java:55
↓ 15 callersMethodadd
(Metric input)
chapter4/src/main/java/com/packtpub/beam/chapter4/ComputeAverage.java:92
↓ 13 callersMethodof
()
util/src/main/java/com/packtpub/beam/util/PrintElements.java:27
↓ 13 callersMethodoutput
()
chapter4/src/main/java/com/packtpub/beam/chapter4/ComputeAverage.java:98
↓ 11 callersMethodof
(String directoryPath)
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:96
↓ 9 callersFunctionmove
(start, direction, speed, duration)
chapter6/src/main/python/utils.py:19
↓ 9 callersMethodof
()
util/src/main/java/com/packtpub/beam/util/MapToLines.java:27
↓ 8 callersMethoddecode
(InputStream inStream)
chapter2/src/main/java/com/packtpub/beam/chapter2/MyStringCoder.java:35
↓ 8 callersFunctionget_expansion_service
(jar="/usr/local/lib/beam-chapter6-1.0.0.jar", args=None)
chapter6/src/main/python/beam_utils.py:5
↓ 7 callersMethodparseFrom
(String tsv)
util/src/main/java/com/packtpub/beam/util/Position.java:29
↓ 7 callersMethodreadAllLines
(InputStream stream)
util/src/main/java/com/packtpub/beam/util/Utils.java:43
↓ 7 callersMethodstart
(self, port)
chapter6/src/main/python/rpc_par_do.py:50
↓ 6 callersMethodof
(Server server)
chapter3/src/main/java/com/packtpub/beam/chapter3/AutoCloseableServer.java:25
↓ 5 callersMethodmain
(String[] args)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDo.java:48
↓ 5 callersMethodtoWords
(String input)
util/src/main/java/com/packtpub/beam/util/Utils.java:37
↓ 4 callersMethodclose
()
chapter3/src/main/java/com/packtpub/beam/chapter3/AutoCloseableServer.java:31
↓ 4 callersMethodcoder
( Coder<K> keyCoder, Coder<J> joinKeyCoder, Coder<V> valueCoder)
chapter4/src/main/java/com/packtpub/beam/chapter4/JoinRow.java:32
↓ 4 callersMethodflush
(self, batch, batchSize, flushTimer, endOfTime)
chapter6/src/main/python/rpc_par_do.py:129
↓ 4 callersMethodnewTempFile
(Path tempDir, int length)
chapter7/src/test/java/com/packtpub/beam/chapter7/StreamingFileReadTest.java:112
↓ 4 callersFunctionrandomPosition
(stamp)
chapter6/src/main/python/test_sport_tracker_motivation.py:14
↓ 4 callersMethodrunRpc
(int port)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDo.java:97
↓ 4 callersMethodts
(String s, Instant stamp)
chapter3/src/test/java/com/packtpub/beam/chapter3/DroppableDataFilter2Test.java:146
↓ 4 callersMethodts
(String s, Instant stamp)
chapter3/src/test/java/com/packtpub/beam/chapter3/DroppableDataFilterTest.java:96
↓ 3 callersFunctionSportTrackerMotivation
(input, shortDuration, longDuration)
chapter6/src/main/python/sport_tracker_motivation.py:143
↓ 3 callersMethodcalculateAverageWordLength
(PCollection<String> words)
chapter2/src/main/java/com/packtpub/beam/chapter2/AverageWordLength.java:76
↓ 3 callersMethodcurrentRestriction
()
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:186
↓ 3 callersMethodencode
(String value, OutputStream outStream)
chapter2/src/main/java/com/packtpub/beam/chapter2/MyStringCoder.java:28
↓ 3 callersMethodflushOutput
( Iterable<ValueWithTimestamp<String>> elements, OutputReceiver<KV<String, Integer>> outputRec
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoStateful.java:304
↓ 3 callersMethodgetLines
(InputStream stream)
util/src/main/java/com/packtpub/beam/util/Utils.java:31
↓ 3 callersMethodifNotNull
(T value)
util/src/main/java/com/packtpub/beam/util/MapToLines.java:52
↓ 3 callersMethodsplitDroppable
( TupleTag<String> mainOutput, TupleTag<String> droppableOutput, PCollection<String> input)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter2.java:130
↓ 3 callersMethodsplitDroppable
( TupleTag<String> mainOutput, TupleTag<String> droppableOutput, PCollection<String> input)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter.java:121
↓ 3 callersMethodtryClaim
(String newFile)
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:177
↓ 2 callersFunctionCalculateAveragePace
(input)
chapter6/src/main/python/sport_tracker_motivation.py:38
↓ 2 callersFunctionComputeLongestWord
(input)
chapter6/src/main/python/max_word_length.py:19
↓ 2 callersFunctionRPCParDoStateful
(input, address="localhost:1234", batchSize=10, maxWaitTime=None)
chapter6/src/main/python/rpc_par_do.py:150
↓ 2 callersFunctionSportTrackerCalc
(input)
chapter6/src/main/python/sport_tracker.py:20
↓ 2 callersMethodapplyRpc
(PCollection<String> input, int port)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDo.java:62
↓ 2 callersMethodapplyRpc
(PCollection<String> input, int port)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoBatch.java:71
↓ 2 callersMethodapplyRpc
(PCollection<String> input, Params params)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoStateful.java:92
↓ 2 callersFunctioncalculateDelta
(deltaLatLon)
chapter6/src/main/python/utils.py:8
↓ 2 callersFunctioncalculateDelta
(deltaLatLon)
chapter6/src/main/python/sport_tracker.py:26
↓ 2 callersMethodcalculateDelta
(double deltaLatLon)
util/src/main/java/com/packtpub/beam/util/Position.java:78
↓ 2 callersMethodclearState
( BagState<ValueWithTimestamp<String>> elements, ValueState<Integer> batchSize, ValueS
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoStateful.java:288
↓ 2 callersMethodcomputeLongestWord
(PCollection<String> lines)
chapter5/src/main/java/com/packtpub/beam/chapter5/SQLMaxWordLength.java:77
↓ 2 callersMethodcomputeLongestWord
(PCollection<String> lines)
chapter2/src/main/java/com/packtpub/beam/chapter2/MaxWordLength.java:71
↓ 2 callersMethodcomputeLongestWordWithTimestamp
( PCollection<String> lines)
chapter2/src/main/java/com/packtpub/beam/chapter2/MaxWordLengthWithTimestamp.java:73
↓ 2 callersMethodcomputeTrackerMetrics
(PCollection<KV<String, Position>> records)
chapter5/src/main/java/com/packtpub/beam/chapter5/SchemaSportTracker.java:130
↓ 2 callersMethodcomputeTrackerMetrics
( PCollection<KV<String, Position>> records)
chapter2/src/main/java/com/packtpub/beam/chapter2/SportTracker.java:126
↓ 2 callersMethodcountWordsInFixedWindows
( PCollection<String> lines, Duration windowLength, int k)
chapter2/src/main/java/com/packtpub/beam/chapter2/TopKWords.java:67
↓ 2 callersMethodcreateFileInDirectory
(Path tempDir, String file)
chapter7/src/test/java/com/packtpub/beam/chapter7/StreamingFileReadTest.java:139
↓ 2 callersFunctiondistance
(p1, p2)
chapter6/src/main/python/utils.py:11
↓ 2 callersMethodflushBuffer
( BagState<KV<Instant, String>> buffer, boolean closed, MultiOutputReceiver output)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter2.java:346
↓ 2 callersMethodflushOutput
( String element, Instant timestamp, boolean closed, MultiOutputReceiver output)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter2.java:335
↓ 2 callersMethodgetResponseFor
(Request request)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCService.java:44
↓ 2 callersMethodhandleProduct
( MapState<K, P> primary, MapState<K, O> other, K primaryKey, J joinKey,
chapter4/src/main/java/com/packtpub/beam/chapter4/StreamingInnerJoin.java:108
↓ 2 callersMethodmapToJoinRow
( SerializableFunction<T, K> keyExtractor, SerializableFunction<T, J> joinKeyExtractor)
chapter4/src/main/java/com/packtpub/beam/chapter4/StreamingInnerJoin.java:266
↓ 2 callersMethodof
(String field)
util/src/main/java/com/packtpub/beam/util/WithStringSchema.java:27
↓ 2 callersMethodofValuesOnly
()
util/src/main/java/com/packtpub/beam/util/MapToLines.java:31
↓ 2 callersMethodparseArgs
(String[] args)
chapter2/src/main/java/com/packtpub/beam/chapter2/TopKWords.java:84
↓ 2 callersMethodregisterCoders
(Pipeline pipeline)
chapter5/src/main/java/com/packtpub/beam/chapter5/SchemaSportTracker.java:211
↓ 2 callersMethodregisterCoders
(Pipeline pipeline)
chapter2/src/main/java/com/packtpub/beam/chapter2/SportTracker.java:191
↓ 2 callersMethodresolveBatch
(self, request, context)
chapter6/src/main/python/rpc_par_do.py:45
↓ 2 callersMethodresolveBatch
( RequestList requestBatch, StreamObserver<ResponseList> responseObserver)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCService.java:33
↓ 2 callersMethodrunWithArgsAndCombiner
(String[] args, TimestampCombiner combiner)
chapter2/src/main/java/com/packtpub/beam/chapter2/MaxWordLengthWithTimestamp.java:50
↓ 2 callersMethodstop
(self)
chapter6/src/main/python/rpc_par_do.py:56
↓ 2 callersMethodstoreResult
( PCollection<String> result, String bootstrapServer, String topic)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter2.java:114
↓ 2 callersMethodstoreResult
( PCollection<String> result, String bootstrapServer, String topic)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter.java:105
↓ 1 callersFunctionComputeBoxedMetrics
(input, duration)
chapter6/src/main/python/sport_tracker_motivation.py:129
↓ 1 callersMethodasPrimary
()
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:153
↓ 1 callersMethodasResidual
()
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:158
↓ 1 callersMethodasStringOfLength
(char character, int length)
chapter7/src/test/java/com/packtpub/beam/chapter7/StreamingFileReadTest.java:131
↓ 1 callersMethodasTestStream
(List<KV<String, Position>> input)
chapter5/src/test/java/com/packtpub/beam/chapter5/SchemaSportTrackerTest.java:115
↓ 1 callersMethodasTestStream
(List<KV<String, Position>> input)
chapter2/src/test/java/com/packtpub/beam/chapter2/SportTrackerTest.java:112
↓ 1 callersMethodcalculateAverageWordLength
(PCollection<String> words)
chapter2/src/main/java/com/packtpub/beam/chapter2/SlidingWindowWordLength.java:64
↓ 1 callersMethodcalculateDistanceOfPositions
(Position first, Position second)
util/src/main/java/com/packtpub/beam/util/Position.java:66
↓ 1 callersMethodcoder
( Coder<K> keyCoder, Coder<J> joinKeyCoder, Coder<L> leftValueCoder, Coder<R> rightValueCoder)
chapter4/src/main/java/com/packtpub/beam/chapter4/JoinResult.java:102
↓ 1 callersFunctioncomputeMetrics
(key, trackPositions)
chapter6/src/main/python/sport_tracker.py:37
↓ 1 callersMethodcomputeMinuteMetrics
( Position first, PriorityQueue<Position> queue, Consumer<TimestampedValue<Metric>> ou
util/src/main/java/com/packtpub/beam/util/ToMetric.java:179
↓ 1 callersMethodcomputeRawMetrics
(Iterable<Row> positions)
chapter5/src/main/java/com/packtpub/beam/chapter5/SchemaSportTracker.java:163
↓ 1 callersMethodcreateInput
(Instant startOfTenSecondWindow)
chapter2/src/test/java/com/packtpub/beam/chapter2/TopKWordsTest.java:52
↓ 1 callersFunctiondistance
(p1, p2)
chapter6/src/main/python/sport_tracker.py:29
↓ 1 callersMethodflushMetrics
(self, flush, key)
chapter6/src/main/python/sport_tracker_motivation.py:107
↓ 1 callersMethodformatMessage
(KV<String, Boolean> message)
util/src/main/java/com/packtpub/beam/util/WriteNotificationsToKafka.java:53
↓ 1 callersMethodgenerateTrack
( int numPointsPerTrack, long timeStart, Random random)
chapter2/src/main/java/com/packtpub/beam/chapter2/SportTracker.java:256
↓ 1 callersMethodgetNewFilesIfAny
( String path, RestrictionTracker<DirectoryWatchRestriction, String> tracker)
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:311
↓ 1 callersMethodgetTimestampPolicyFactory
()
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter2.java:106
↓ 1 callersMethodgetTimestampPolicyFactory
()
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter.java:97
↓ 1 callersMethodloadDataFromStream
(InputStream stream)
chapter5/src/test/java/com/packtpub/beam/chapter5/SchemaSportTrackerTest.java:126
↓ 1 callersMethodloadDataFromStream
(InputStream stream)
chapter2/src/test/java/com/packtpub/beam/chapter2/SportTrackerTest.java:123
↓ 1 callersMethodmerge
(MetricAccumulator acc)
chapter4/src/main/java/com/packtpub/beam/chapter4/ComputeAverage.java:102
next →1–100 of 408, ranked by callers