Side input patterns
The samples on this page show you common Beam side input patterns. A side input is an additional input that your DoFn can access each time it processes an element in the input PCollection. For more information, see the programming guide section on side inputs.
If you are trying to enrich your data by doing a key-value lookup to a remote service, you may first want to consider the Enrichment transform which can abstract away some of the details of side inputs and provide additional benefits like client-side throttling.
- Java SDK
- Python SDK
Slowly updating global window side inputs
You can retrieve side inputs from global windows to use them in a pipeline job with non-global windows, like a FixedWindow.
To slowly update global window side inputs in pipelines with non-global windows:
Write a
DoFnthat periodically pulls data from a bounded source into a global window.a. Use the
GenerateSequencesource transform to periodically emit a value.b. Instantiate a data-driven trigger that activates on each element and pulls data from a bounded source.
c. Fire the trigger to pass the data into the global window.
Create the side input for downstream transforms. The side input should fit into memory.
The global window side input triggers on processing time, so the main pipeline non-deterministically matches the side input to elements in event time.
For instance, the following code sample uses a Map to create a DoFn. The Map becomes a View.asSingleton side input thatâs rebuilt on each counter tick. The side input updates every 5 seconds in order to demonstrate the workflow. In a real-world scenario, the side input would typically update every few hours or once per day.
public static void sideInputPatterns() {
// This pipeline uses View.asSingleton for a placeholder external service.
// Run in debug mode to see the output.
Pipeline p = Pipeline.create();
// Create a side input that updates every 5 seconds.
// View as an iterable, not singleton, so that if we happen to trigger more
// than once before Latest.globally is computed we can handle both elements.
PCollectionView<Iterable<Map<String, String>>> mapIterable =
p.apply(GenerateSequence.from(0).withRate(1, Duration.standardSeconds(5L)))
.apply(
ParDo.of(
new DoFn<Long, Map<String, String>>() {
@ProcessElement
public void process(
@Element Long input,
@Timestamp Instant timestamp,
OutputReceiver<Map<String, String>> o) {
// Replace map with test data from the placeholder external service.
// Add external reads here.
o.output(PlaceholderExternalService.readTestData(timestamp));
}
}))
.apply(
Window.<Map<String, String>>into(new GlobalWindows())
.triggering(Repeatedly.forever(AfterProcessingTime.pastFirstElementInPane()))
.discardingFiredPanes())
.apply(Latest.globally())
.apply(View.asIterable());
// Consume side input. GenerateSequence generates test data.
// Use a real source (like PubSubIO or KafkaIO) in production.
p.apply(GenerateSequence.from(0).withRate(1, Duration.standardSeconds(1L)))
.apply(Window.into(FixedWindows.of(Duration.standardSeconds(1))))
.apply(Sum.longsGlobally().withoutDefaults())
.apply(
ParDo.of(
new DoFn<Long, KV<Long, Long>>() {
@ProcessElement
public void process(
@Timestamp Instant timestamp,
@Element Long element,
@SideInput("mapIterable") Iterable<Map<String, String>> si,
OutputReceiver<KV<Long, Long>> receiver) {
// Take an element from the side input iterable (likely length 1)
Map<String, String> keyMap = si.iterator().next();
receiver.outputWithTimestamp(KV.of(1L, element), Instant.now());
LOG.info(
"Value is {} with timestamp {}, using key A from side input with time {}.",
element,
timestamp.toString(DateTimeFormat.forPattern("HH:mm:ss")),
keyMap.get("Key_A"));
}
})
.withSideInput("mapIterable", mapIterable));
p.run();
}
/** Placeholder class that represents an external service generating test data. */
public static class PlaceholderExternalService {
public static Map<String, String> readTestData(Instant timestamp) {
Map<String, String> map = new HashMap<>();
map.put("Key_A", timestamp.toString(DateTimeFormat.forPattern("HH:mm:ss")));
return map;
}
}The Python sample uses PeriodicImpulse to re-read the placeholder external service on a fixed interval, and re-publishes the result into the global window on every firing. Use Latest.Globally().without_defaults() rather than Latest.Globally(): the variant with defaults adds its own side input, which stops the transform from emitting more than once. For more information, see Issue 35934.
from apache_beam.transforms import combiners
from apache_beam.transforms import trigger
from apache_beam.transforms import window
from apache_beam.transforms.periodicsequence import PeriodicImpulse
# To run indefinitely, pass MAX_TIMESTAMP from apache_beam.utils.timestamp
# as last_timestamp.
# Placeholder that represents an external service, such as a database or a
# configuration endpoint. Replace it with the external read of your choice.
def read_from_placeholder_external_service(refresh_timestamp):
return {'Key_A': str(refresh_timestamp)}
def enrich_with_side_input(element, config):
# The side input is read as an iterable rather than as a singleton,
# because the global window side input can hold more than one element if
# it fires again before the value is consumed.
latest_config = next(iter(config), {})
return element, latest_config.get('Key_A')
# Create pipeline.
pipeline = beam.Pipeline()
# Periodically read the external data into the global window.
# Repeatedly(AfterCount(1)) emits a new pane for every impulse, and the
# DISCARDING accumulation mode drops the previous value, so that
# Latest.Globally only considers the most recent element.
# Use Latest.Globally().without_defaults(): the variant with defaults adds
# its own side input, which stops the transform from emitting more than once.
side_input = (
pipeline
| 'SideInputImpulse' >> PeriodicImpulse(
first_timestamp, last_timestamp, side_input_interval)
| 'ReadExternalData' >> beam.Map(read_from_placeholder_external_service)
| 'WindowSideInput' >> beam.WindowInto(
window.GlobalWindows(),
trigger=trigger.Repeatedly(trigger.AfterCount(1)),
accumulation_mode=trigger.AccumulationMode.DISCARDING)
| 'GetLatest' >> combiners.Latest.Globally().without_defaults())
# Consume the side input from a main input that uses non-global windows.
# PeriodicImpulse generates test data. Use a real streaming source, such as
# PubSubIO or KafkaIO, in production.
result = (
pipeline
| 'MainInputImpulse' >> PeriodicImpulse(
first_timestamp,
last_timestamp,
main_input_interval,
apply_windowing=True)
| 'ApplySideInput' >> beam.Map(
enrich_with_side_input, config=beam.pvalue.AsIter(side_input)))Slowly updating side input using windowing
You can read side input data periodically into distinct PCollection windows. When you apply the side input to your main input, each main input window is automatically matched to a single side input window. This guarantees consistency on the duration of the single window, meaning that each window on the main input will be matched to a single version of side input data.
To read side input data periodically into distinct PCollection windows:
- Use the PeriodicImpulse or PeriodicSequence PTransform to:
- Generate an infinite sequence of elements at required processing time intervals
- Assign them to separate windows.
- Fetch data using SDF Read or ReadAll PTransform triggered by arrival of PCollection element.
- Apply the side input.
PCollectionView<List<Long>> sideInput =
p.apply(
"SIImpulse",
PeriodicImpulse.create()
.startAt(startAt)
.stopAt(stopAt)
.withInterval(interval1)
.applyWindowing())
.apply(
"FileToRead",
ParDo.of(
new DoFn<Instant, String>() {
@DoFn.ProcessElement
public void process(@Element Instant notUsed, OutputReceiver<String> o) {
o.output(fileToRead);
}
}))
.apply(FileIO.matchAll())
.apply(FileIO.readMatches())
.apply(TextIO.readFiles())
.apply(
ParDo.of(
new DoFn<String, String>() {
@ProcessElement
public void process(@Element String src, OutputReceiver<String> o) {
o.output(src);
}
}))
.apply(Combine.globally(Count.<String>combineFn()).withoutDefaults())
.apply(View.asList());
PCollection<Instant> mainInput =
p.apply(
"MIImpulse",
PeriodicImpulse.create()
.startAt(startAt.minus(Duration.standardSeconds(1)))
.stopAt(stopAt.minus(Duration.standardSeconds(1)))
.withInterval(interval2)
.applyWindowing());
// Consume side input. GenerateSequence generates test data.
// Use a real source (like PubSubIO or KafkaIO) in production.
PCollection<Long> result =
mainInput.apply(
"generateOutput",
ParDo.of(
new DoFn<Instant, Long>() {
@ProcessElement
public void process(
@SideInput("sideInput") List<Long> sideInputValue,
OutputReceiver<Long> receiver) {
receiver.output((long) sideInputValue.size());
}
})
.withSideInput("sideInput", sideInput));from apache_beam.transforms.periodicsequence import PeriodicImpulse
from apache_beam.transforms.window import TimestampedValue
from apache_beam.transforms import window
# from apache_beam.utils.timestamp import MAX_TIMESTAMP
# last_timestamp = MAX_TIMESTAMP to go on indefninitely
# Any user-defined function.
# cross join is used as an example.
def cross_join(left, rights):
for x in rights:
yield (left, x)
# Create pipeline.
pipeline = beam.Pipeline()
side_input = (
pipeline
| 'PeriodicImpulse' >> PeriodicImpulse(
first_timestamp, last_timestamp, interval, True)
| 'MapToFileName' >> beam.Map(lambda x: src_file_pattern + str(x))
| 'ReadFromFile' >> beam.io.ReadAllFromText())
main_input = (
pipeline
| 'MpImpulse' >> beam.Create(sample_main_input_elements)
|
'MapMpToTimestamped' >> beam.Map(lambda src: TimestampedValue(src, src))
| 'WindowMpInto' >> beam.WindowInto(
window.FixedWindows(main_input_windowing_interval)))
result = (
main_input
| 'ApplyCrossJoin' >> beam.FlatMap(
cross_join, rights=beam.pvalue.AsIter(side_input)))Last updated on 2026/10/02
Have you found everything you were looking for?
Was it all useful and clear? Is there anything that you would like to change? Let us know!

