Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/PacktPublishing/Building-Big-Data-Pipelines-with-Apache-Beam
/ functions
Functions
408 in github.com/PacktPublishing/Building-Big-Data-Pipelines-with-Apache-Beam
⨍
Functions
408
◇
Types & classes
144
↳
Endpoints
1
↓ 367 callers
Method
apply
()
chapter7/src/main/java/com/packtpub/beam/chapter7/PiSampler.java:45
↓ 226 callers
Method
of
( J joinKey, @Nullable K leftKey, @Nullable L leftValue, @Nullable K rightKey, @
chapter4/src/main/java/com/packtpub/beam/chapter4/JoinResult.java:31
↓ 165 callers
Method
of
(long numSamples, long parallelism)
chapter7/src/main/java/com/packtpub/beam/chapter7/PiSampler.java:54
↓ 123 callers
Method
create
(String name)
util/src/main/java/com/packtpub/beam/util/DebugOutput.java:27
↓ 62 callers
Method
of
(Coder<T> valueCoder)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoStateful.java:154
↓ 49 callers
Method
of
()
util/src/main/java/com/packtpub/beam/util/Tokenize.java:25
↓ 34 callers
Method
move
( double latitudeDirection, double longitudeDirection, double speedMeterPerSec, long t
util/src/main/java/com/packtpub/beam/util/Position.java:50
↓ 28 callers
Method
random
(long stamp)
util/src/main/java/com/packtpub/beam/util/Position.java:42
↓ 27 callers
Method
distance
(Position other)
util/src/main/java/com/packtpub/beam/util/Position.java:46
↓ 19 callers
Method
collect
( J joinKey, K primaryKey, Object primaryValue, K otherKey,
chapter4/src/main/java/com/packtpub/beam/chapter4/StreamingInnerJoin.java:55
↓ 15 callers
Method
add
(Metric input)
chapter4/src/main/java/com/packtpub/beam/chapter4/ComputeAverage.java:92
↓ 13 callers
Method
of
()
util/src/main/java/com/packtpub/beam/util/PrintElements.java:27
↓ 13 callers
Method
output
()
chapter4/src/main/java/com/packtpub/beam/chapter4/ComputeAverage.java:98
↓ 11 callers
Method
of
(String directoryPath)
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:96
↓ 9 callers
Function
move
(start, direction, speed, duration)
chapter6/src/main/python/utils.py:19
↓ 9 callers
Method
of
()
util/src/main/java/com/packtpub/beam/util/MapToLines.java:27
↓ 8 callers
Method
decode
(InputStream inStream)
chapter2/src/main/java/com/packtpub/beam/chapter2/MyStringCoder.java:35
↓ 8 callers
Function
get_expansion_service
(jar="/usr/local/lib/beam-chapter6-1.0.0.jar", args=None)
chapter6/src/main/python/beam_utils.py:5
↓ 7 callers
Method
parseFrom
(String tsv)
util/src/main/java/com/packtpub/beam/util/Position.java:29
↓ 7 callers
Method
readAllLines
(InputStream stream)
util/src/main/java/com/packtpub/beam/util/Utils.java:43
↓ 7 callers
Method
start
(self, port)
chapter6/src/main/python/rpc_par_do.py:50
↓ 6 callers
Method
of
(Server server)
chapter3/src/main/java/com/packtpub/beam/chapter3/AutoCloseableServer.java:25
↓ 5 callers
Method
main
(String[] args)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDo.java:48
↓ 5 callers
Method
toWords
(String input)
util/src/main/java/com/packtpub/beam/util/Utils.java:37
↓ 4 callers
Method
close
()
chapter3/src/main/java/com/packtpub/beam/chapter3/AutoCloseableServer.java:31
↓ 4 callers
Method
coder
( Coder<K> keyCoder, Coder<J> joinKeyCoder, Coder<V> valueCoder)
chapter4/src/main/java/com/packtpub/beam/chapter4/JoinRow.java:32
↓ 4 callers
Method
flush
(self, batch, batchSize, flushTimer, endOfTime)
chapter6/src/main/python/rpc_par_do.py:129
↓ 4 callers
Method
newTempFile
(Path tempDir, int length)
chapter7/src/test/java/com/packtpub/beam/chapter7/StreamingFileReadTest.java:112
↓ 4 callers
Function
randomPosition
(stamp)
chapter6/src/main/python/test_sport_tracker_motivation.py:14
↓ 4 callers
Method
runRpc
(int port)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDo.java:97
↓ 4 callers
Method
ts
(String s, Instant stamp)
chapter3/src/test/java/com/packtpub/beam/chapter3/DroppableDataFilter2Test.java:146
↓ 4 callers
Method
ts
(String s, Instant stamp)
chapter3/src/test/java/com/packtpub/beam/chapter3/DroppableDataFilterTest.java:96
↓ 3 callers
Function
SportTrackerMotivation
(input, shortDuration, longDuration)
chapter6/src/main/python/sport_tracker_motivation.py:143
↓ 3 callers
Method
calculateAverageWordLength
(PCollection<String> words)
chapter2/src/main/java/com/packtpub/beam/chapter2/AverageWordLength.java:76
↓ 3 callers
Method
currentRestriction
()
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:186
↓ 3 callers
Method
encode
(String value, OutputStream outStream)
chapter2/src/main/java/com/packtpub/beam/chapter2/MyStringCoder.java:28
↓ 3 callers
Method
flushOutput
( Iterable<ValueWithTimestamp<String>> elements, OutputReceiver<KV<String, Integer>> outputRec
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoStateful.java:304
↓ 3 callers
Method
getLines
(InputStream stream)
util/src/main/java/com/packtpub/beam/util/Utils.java:31
↓ 3 callers
Method
ifNotNull
(T value)
util/src/main/java/com/packtpub/beam/util/MapToLines.java:52
↓ 3 callers
Method
splitDroppable
( TupleTag<String> mainOutput, TupleTag<String> droppableOutput, PCollection<String> input)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter2.java:130
↓ 3 callers
Method
splitDroppable
( TupleTag<String> mainOutput, TupleTag<String> droppableOutput, PCollection<String> input)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter.java:121
↓ 3 callers
Method
tryClaim
(String newFile)
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:177
↓ 2 callers
Function
CalculateAveragePace
(input)
chapter6/src/main/python/sport_tracker_motivation.py:38
↓ 2 callers
Function
ComputeLongestWord
(input)
chapter6/src/main/python/max_word_length.py:19
↓ 2 callers
Function
RPCParDoStateful
(input, address="localhost:1234", batchSize=10, maxWaitTime=None)
chapter6/src/main/python/rpc_par_do.py:150
↓ 2 callers
Function
SportTrackerCalc
(input)
chapter6/src/main/python/sport_tracker.py:20
↓ 2 callers
Method
applyRpc
(PCollection<String> input, int port)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDo.java:62
↓ 2 callers
Method
applyRpc
(PCollection<String> input, int port)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoBatch.java:71
↓ 2 callers
Method
applyRpc
(PCollection<String> input, Params params)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoStateful.java:92
↓ 2 callers
Function
calculateDelta
(deltaLatLon)
chapter6/src/main/python/utils.py:8
↓ 2 callers
Function
calculateDelta
(deltaLatLon)
chapter6/src/main/python/sport_tracker.py:26
↓ 2 callers
Method
calculateDelta
(double deltaLatLon)
util/src/main/java/com/packtpub/beam/util/Position.java:78
↓ 2 callers
Method
clearState
( BagState<ValueWithTimestamp<String>> elements, ValueState<Integer> batchSize, ValueS
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCParDoStateful.java:288
↓ 2 callers
Method
computeLongestWord
(PCollection<String> lines)
chapter5/src/main/java/com/packtpub/beam/chapter5/SQLMaxWordLength.java:77
↓ 2 callers
Method
computeLongestWord
(PCollection<String> lines)
chapter2/src/main/java/com/packtpub/beam/chapter2/MaxWordLength.java:71
↓ 2 callers
Method
computeLongestWordWithTimestamp
( PCollection<String> lines)
chapter2/src/main/java/com/packtpub/beam/chapter2/MaxWordLengthWithTimestamp.java:73
↓ 2 callers
Method
computeTrackerMetrics
(PCollection<KV<String, Position>> records)
chapter5/src/main/java/com/packtpub/beam/chapter5/SchemaSportTracker.java:130
↓ 2 callers
Method
computeTrackerMetrics
( PCollection<KV<String, Position>> records)
chapter2/src/main/java/com/packtpub/beam/chapter2/SportTracker.java:126
↓ 2 callers
Method
countWordsInFixedWindows
( PCollection<String> lines, Duration windowLength, int k)
chapter2/src/main/java/com/packtpub/beam/chapter2/TopKWords.java:67
↓ 2 callers
Method
createFileInDirectory
(Path tempDir, String file)
chapter7/src/test/java/com/packtpub/beam/chapter7/StreamingFileReadTest.java:139
↓ 2 callers
Function
distance
(p1, p2)
chapter6/src/main/python/utils.py:11
↓ 2 callers
Method
flushBuffer
( BagState<KV<Instant, String>> buffer, boolean closed, MultiOutputReceiver output)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter2.java:346
↓ 2 callers
Method
flushOutput
( String element, Instant timestamp, boolean closed, MultiOutputReceiver output)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter2.java:335
↓ 2 callers
Method
getResponseFor
(Request request)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCService.java:44
↓ 2 callers
Method
handleProduct
( MapState<K, P> primary, MapState<K, O> other, K primaryKey, J joinKey,
chapter4/src/main/java/com/packtpub/beam/chapter4/StreamingInnerJoin.java:108
↓ 2 callers
Method
mapToJoinRow
( SerializableFunction<T, K> keyExtractor, SerializableFunction<T, J> joinKeyExtractor)
chapter4/src/main/java/com/packtpub/beam/chapter4/StreamingInnerJoin.java:266
↓ 2 callers
Method
of
(String field)
util/src/main/java/com/packtpub/beam/util/WithStringSchema.java:27
↓ 2 callers
Method
ofValuesOnly
()
util/src/main/java/com/packtpub/beam/util/MapToLines.java:31
↓ 2 callers
Method
parseArgs
(String[] args)
chapter2/src/main/java/com/packtpub/beam/chapter2/TopKWords.java:84
↓ 2 callers
Method
registerCoders
(Pipeline pipeline)
chapter5/src/main/java/com/packtpub/beam/chapter5/SchemaSportTracker.java:211
↓ 2 callers
Method
registerCoders
(Pipeline pipeline)
chapter2/src/main/java/com/packtpub/beam/chapter2/SportTracker.java:191
↓ 2 callers
Method
resolveBatch
(self, request, context)
chapter6/src/main/python/rpc_par_do.py:45
↓ 2 callers
Method
resolveBatch
( RequestList requestBatch, StreamObserver<ResponseList> responseObserver)
chapter3/src/main/java/com/packtpub/beam/chapter3/RPCService.java:33
↓ 2 callers
Method
runWithArgsAndCombiner
(String[] args, TimestampCombiner combiner)
chapter2/src/main/java/com/packtpub/beam/chapter2/MaxWordLengthWithTimestamp.java:50
↓ 2 callers
Method
stop
(self)
chapter6/src/main/python/rpc_par_do.py:56
↓ 2 callers
Method
storeResult
( PCollection<String> result, String bootstrapServer, String topic)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter2.java:114
↓ 2 callers
Method
storeResult
( PCollection<String> result, String bootstrapServer, String topic)
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter.java:105
↓ 1 callers
Function
ComputeBoxedMetrics
(input, duration)
chapter6/src/main/python/sport_tracker_motivation.py:129
↓ 1 callers
Method
asPrimary
()
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:153
↓ 1 callers
Method
asResidual
()
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:158
↓ 1 callers
Method
asStringOfLength
(char character, int length)
chapter7/src/test/java/com/packtpub/beam/chapter7/StreamingFileReadTest.java:131
↓ 1 callers
Method
asTestStream
(List<KV<String, Position>> input)
chapter5/src/test/java/com/packtpub/beam/chapter5/SchemaSportTrackerTest.java:115
↓ 1 callers
Method
asTestStream
(List<KV<String, Position>> input)
chapter2/src/test/java/com/packtpub/beam/chapter2/SportTrackerTest.java:112
↓ 1 callers
Method
calculateAverageWordLength
(PCollection<String> words)
chapter2/src/main/java/com/packtpub/beam/chapter2/SlidingWindowWordLength.java:64
↓ 1 callers
Method
calculateDistanceOfPositions
(Position first, Position second)
util/src/main/java/com/packtpub/beam/util/Position.java:66
↓ 1 callers
Method
coder
( Coder<K> keyCoder, Coder<J> joinKeyCoder, Coder<L> leftValueCoder, Coder<R> rightValueCoder)
chapter4/src/main/java/com/packtpub/beam/chapter4/JoinResult.java:102
↓ 1 callers
Function
computeMetrics
(key, trackPositions)
chapter6/src/main/python/sport_tracker.py:37
↓ 1 callers
Method
computeMinuteMetrics
( Position first, PriorityQueue<Position> queue, Consumer<TimestampedValue<Metric>> ou
util/src/main/java/com/packtpub/beam/util/ToMetric.java:179
↓ 1 callers
Method
computeRawMetrics
(Iterable<Row> positions)
chapter5/src/main/java/com/packtpub/beam/chapter5/SchemaSportTracker.java:163
↓ 1 callers
Method
createInput
(Instant startOfTenSecondWindow)
chapter2/src/test/java/com/packtpub/beam/chapter2/TopKWordsTest.java:52
↓ 1 callers
Function
distance
(p1, p2)
chapter6/src/main/python/sport_tracker.py:29
↓ 1 callers
Method
flushMetrics
(self, flush, key)
chapter6/src/main/python/sport_tracker_motivation.py:107
↓ 1 callers
Method
formatMessage
(KV<String, Boolean> message)
util/src/main/java/com/packtpub/beam/util/WriteNotificationsToKafka.java:53
↓ 1 callers
Method
generateTrack
( int numPointsPerTrack, long timeStart, Random random)
chapter2/src/main/java/com/packtpub/beam/chapter2/SportTracker.java:256
↓ 1 callers
Method
getNewFilesIfAny
( String path, RestrictionTracker<DirectoryWatchRestriction, String> tracker)
chapter7/src/main/java/com/packtpub/beam/chapter7/StreamingFileRead.java:311
↓ 1 callers
Method
getTimestampPolicyFactory
()
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter2.java:106
↓ 1 callers
Method
getTimestampPolicyFactory
()
chapter3/src/main/java/com/packtpub/beam/chapter3/DroppableDataFilter.java:97
↓ 1 callers
Method
loadDataFromStream
(InputStream stream)
chapter5/src/test/java/com/packtpub/beam/chapter5/SchemaSportTrackerTest.java:126
↓ 1 callers
Method
loadDataFromStream
(InputStream stream)
chapter2/src/test/java/com/packtpub/beam/chapter2/SportTrackerTest.java:123
↓ 1 callers
Method
merge
(MetricAccumulator acc)
chapter4/src/main/java/com/packtpub/beam/chapter4/ComputeAverage.java:102
next →
1–100 of 408, ranked by callers