diff --git a/pyflink-walkthrough/payment_msg_proccessing.py b/pyflink-walkthrough/payment_msg_proccessing.py index bc2a1c7..f486f1a 100644 --- a/pyflink-walkthrough/payment_msg_proccessing.py +++ b/pyflink-walkthrough/payment_msg_proccessing.py @@ -76,7 +76,7 @@ def log_processing(): t_env.from_path("payment_msg") \ .select(call('province_id_to_name', col('provinceId')).alias("province"), col('payAmount')) \ .group_by(col('province')) \ - .select(col('province'), call('sum', col('payAmount').alias("pay_amount"))) \ + .select(col('province'), call('sum', col('payAmount')).alias("pay_amount")) \ .execute_insert("es_sink")