diff --git a/bootcamp/materials/4-apache-flink-training/src/job/start_job.py b/bootcamp/materials/4-apache-flink-training/src/job/start_job.py index 3a37d2b2..5f66732e 100644 --- a/bootcamp/materials/4-apache-flink-training/src/job/start_job.py +++ b/bootcamp/materials/4-apache-flink-training/src/job/start_job.py @@ -4,9 +4,8 @@ import json import requests from pyflink.table import EnvironmentSettings, DataTypes, TableEnvironment, StreamTableEnvironment - def create_processed_events_sink_kafka(t_env): - table_name = "process_events_kafka" + table_name = "raw_events_kafka" kafka_key = os.environ.get("KAFKA_WEB_TRAFFIC_KEY", "") kafka_secret = os.environ.get("KAFKA_WEB_TRAFFIC_SECRET", "") sasl_config = f'org.apache.kafka.common.security.plain.PlainLoginModule required username="{kafka_key}" password="{kafka_secret}";'