in pipelines/marketing_intelligence/marketing_intelligence_pipeline/pipeline.py [0:0]
def _extract(p: Pipeline, subscription: str) -> PCollection[str]:
msgs: PCollection[bytes] = p | "Read subscription" >> beam.io.ReadFromPubSub(
subscription=subscription)
return msgs | "Parse and format Input" >> beam.Map(_format_input)