apache/beam · error · RuntimeError

Invalid initial position in stream

Error message

Invalid initial position in stream: {}

What it means

InitialPositionInStream.validate_param checks that the given position name is an attribute of the InitialPositionInStream class (TRIM_HORIZON or AT_TIMESTAMP). An unknown, misspelled, or empty-but-truthy string raises RuntimeError. It is a validation helper for the Kinesis ReadFromKinesis source.

Solutions

  1. Use InitialPositionInStream.TRIM_HORIZON or InitialPositionInStream.AT_TIMESTAMP constants.
  2. If specifying AT_TIMESTAMP, also pass a valid timestamp (at_timestamp) argument.
  3. Normalize/validate externally-provided config values against the class attributes before constructing the source.

Example fix

# before
ReadFromKinesis(stream_name, initial_position_in_stream='trim_horizon')
# after
ReadFromKinesis(stream_name, initial_position_in_stream=InitialPositionInStream.TRIM_HORIZON)
Defensive patterns

Strategy: validation

Validate before calling

if param and not hasattr(InitialPositionInStream, param):
    raise ValueError('position must be TRIM_HORIZON or AT_TIMESTAMP')

Type guard

def is_valid_initial_position(param) -> bool:
    return not param or hasattr(InitialPositionInStream, param)

Try / catch

try:
    InitialPositionInStream.validate_param(position)
except RuntimeError as e:
    log.error('Bad initial position %r; use InitialPositionInStream constants', position)

Prevention

When it happens

Trigger: ReadFromKinesis(..., initial_position_in_stream='trim_horizon') or 'LATEST' or any string that is not exactly 'TRIM_HORIZON' or 'AT_TIMESTAMP'; also any truthy invalid value passed through the static validate_param.

Common situations: Lowercasing the position name; using Kafka's 'latest'/'earliest' vocabulary with Kinesis; dynamically building the parameter from config files.

Understand the failure class

Background: Invalid enum value errors: "Unknown type", "Invalid scope", "must be one of" — when a string is not on the library's allowed list — this error's family across 23 libraries.

Related errors


AI-assisted analysis of apache/beam@12126d8942 (2026-09-13). Data as JSON: /api/errors/dee85cd033b94704. Report an issue: GitHub.

Appendix: source

Thrown at sdks/python/apache_beam/io/kinesis.py:337

                max_capacity_per_shard=max_capacity_per_shard,
                watermark_policy=watermark_policy,
                watermark_idle_duration_threshold=
                watermark_idle_duration_threshold,
                rate_limit=rate_limit,
            )),
        expansion_service or default_io_expansion_service(),
    )


class InitialPositionInStream:
  LATEST = 'LATEST'
  TRIM_HORIZON = 'TRIM_HORIZON'
  AT_TIMESTAMP = 'AT_TIMESTAMP'

  @staticmethod
  def validate_param(param):
    if param and not hasattr(InitialPositionInStream, param):
      raise RuntimeError('Invalid initial position in stream: {}'.format(param))


class WatermarkPolicy:
  PROCESSING_TYPE = 'PROCESSING_TYPE'
  ARRIVAL_TIME = 'ARRIVAL_TIME'

  @staticmethod
  def validate_param(param):
    if param and not hasattr(WatermarkPolicy, param):
      raise RuntimeError('Invalid watermark policy: {}'.format(param))

View on GitHub (pinned to 12126d8942)