{"record":{"id":"dee85cd033b94704","repo":"apache/beam","slug":"invalid-initial-position-in-stream","errorCode":null,"errorMessage":"Invalid initial position in stream: {}","messagePattern":"Invalid initial position in stream: (.+?)","errorType":"validation","errorClass":"RuntimeError","httpStatus":null,"severity":"error","filePath":"sdks/python/apache_beam/io/kinesis.py","lineNumber":337,"sourceCode":"                max_capacity_per_shard=max_capacity_per_shard,\n                watermark_policy=watermark_policy,\n                watermark_idle_duration_threshold=\n                watermark_idle_duration_threshold,\n                rate_limit=rate_limit,\n            )),\n        expansion_service or default_io_expansion_service(),\n    )\n\n\nclass InitialPositionInStream:\n  LATEST = 'LATEST'\n  TRIM_HORIZON = 'TRIM_HORIZON'\n  AT_TIMESTAMP = 'AT_TIMESTAMP'\n\n  @staticmethod\n  def validate_param(param):\n    if param and not hasattr(InitialPositionInStream, param):\n      raise RuntimeError('Invalid initial position in stream: {}'.format(param))\n\n\nclass WatermarkPolicy:\n  PROCESSING_TYPE = 'PROCESSING_TYPE'\n  ARRIVAL_TIME = 'ARRIVAL_TIME'\n\n  @staticmethod\n  def validate_param(param):\n    if param and not hasattr(WatermarkPolicy, param):\n      raise RuntimeError('Invalid watermark policy: {}'.format(param))\n","sourceCodeStart":319,"sourceCodeEnd":348,"githubUrl":"https://github.com/apache/beam/blob/12126d8942aaf848030c478b4c6a28c6af861c66/sdks/python/apache_beam/io/kinesis.py#L319-L348","documentation":"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.","triggerScenarios":"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.","commonSituations":"Lowercasing the position name; using Kafka's 'latest'/'earliest' vocabulary with Kinesis; dynamically building the parameter from config files.","solutions":["Use InitialPositionInStream.TRIM_HORIZON or InitialPositionInStream.AT_TIMESTAMP constants.","If specifying AT_TIMESTAMP, also pass a valid timestamp (at_timestamp) argument.","Normalize/validate externally-provided config values against the class attributes before constructing the source."],"exampleFix":"# before\nReadFromKinesis(stream_name, initial_position_in_stream='trim_horizon')\n# after\nReadFromKinesis(stream_name, initial_position_in_stream=InitialPositionInStream.TRIM_HORIZON)","handlingStrategy":"validation","validationCode":"if param and not hasattr(InitialPositionInStream, param):\n    raise ValueError('position must be TRIM_HORIZON or AT_TIMESTAMP')","typeGuard":"def is_valid_initial_position(param) -> bool:\n    return not param or hasattr(InitialPositionInStream, param)","tryCatchPattern":"try:\n    InitialPositionInStream.validate_param(position)\nexcept RuntimeError as e:\n    log.error('Bad initial position %r; use InitialPositionInStream constants', position)","preventionTips":["Use InitialPositionInStream.TRIM_HORIZON / AT_TIMESTAMP constants","Do not reuse Kafka 'latest'/'earliest' terms for Kinesis","Normalize user/config input to uppercase enum names before validation"],"tags":["apache-beam","python","kinesis","invalid-enum-value","validation"],"backgroundTag":"invalid-enum-value","analyzedSha":"12126d8942aaf848030c478b4c6a28c6af861c66","analyzedAt":"2026-09-13T01:50:10.254Z","contentChangedAt":"2026-09-13T01:50:10.254Z","schemaVersion":2},"datasetVersion":"2026-09-20T03:17:13.778Z"}