{"record":{"id":"db733f3baa07fb93","repo":"apache/beam","slug":"invalid-watermark-policy","errorCode":null,"errorMessage":"Invalid watermark policy: {}","messagePattern":"Invalid watermark policy: (.+?)","errorType":"validation","errorClass":"RuntimeError","httpStatus":null,"severity":"error","filePath":"sdks/python/apache_beam/io/kinesis.py","lineNumber":347,"sourceCode":"class 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":329,"sourceCodeEnd":348,"githubUrl":"https://github.com/apache/beam/blob/12126d8942aaf848030c478b4c6a28c6af861c66/sdks/python/apache_beam/io/kinesis.py#L329-L348","documentation":"WatermarkPolicy.validate_param checks that a given watermark policy name is a valid class attribute of WatermarkPolicy (e.g. 'ARRIVAL_TIME' or 'PROCESSING_TYPE'). When a truthy string is passed that does not match any defined policy, a RuntimeError is raised with the offending value. This fails fast at pipeline construction time instead of producing a confusing failure inside the Kinesis reader later.","triggerScenarios":"Calling ReadDataFromKinesis/beam.io.kinesis.ReadFromKinesis with watermark_policy set to any string not present as an attribute on WatermarkPolicy, e.g. watermark_policy='Arrival_Time', 'arrival_time' or a typo like 'ARRIVALL_TIME'. Empty/None passes validation because `if param` short-circuits.","commonSituations":"Hand-typing the policy name in pipeline options or YAML templates instead of referencing WatermarkPolicy.ARRIVAL_TIME / WatermarkPolicy.PROCESSING_TYPE; case-sensitivity mistakes after copying config between jobs; renaming/upgrade where the allowed policy values changed.","solutions":["Use one of the defined constants: apache_beam.io.kinesis.WatermarkPolicy.ARRIVAL_TIME or WatermarkPolicy.PROCESSING_TYPE, instead of a raw string.","Check the exact spelling and case of the policy string; the check is hasattr-based and case-sensitive.","Pass None or omit the parameter if you do not need watermark updates, since empty values are accepted.","Inspect the WatermarkPolicy class in your installed Beam version (sdks/python/apache_beam/io/kinesis.py) to see which values are supported."],"exampleFix":"// before\nReadDataFromKinesis(input_stream='stream', watermark_policy='Arrival_Time')\n// after\nfrom apache_beam.io.kinesis import WatermarkPolicy, ReadDataFromKinesis\nReadDataFromKinesis(input_stream='stream', watermark_policy=WatermarkPolicy.ARRIVAL_TIME)","handlingStrategy":"validation","validationCode":"from apache_beam.io.kinesis import WatermarkPolicy\n\ndef valid_watermark_policy(p):\n    return not p or hasattr(WatermarkPolicy, p)\n\nassert valid_watermark_policy(watermark_policy), f\"unknown policy: {watermark_policy}\"","typeGuard":"def is_watermark_policy(value) -> bool:\n    from apache_beam.io.kinesis import WatermarkPolicy\n    return isinstance(value, str) and hasattr(WatermarkPolicy, value)","tryCatchPattern":null,"preventionTips":["Always reference WatermarkPolicy.ARRIVAL_TIME / WatermarkPolicy.PROCESSING_TYPE constants instead of raw strings","Remember validation is case-sensitive; avoid hand-typing policy names","Validate pipeline options early in a unit test that constructs the transform"],"tags":["python","apache-beam","kinesis","invalid-enum-value","config-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"}