DataFlow中DoFn类send_email函数无输出问题排查求助
DataFlow中send_email函数在DataflowRunner下失效问题
问题详情
- 数据掩码、BigQuery数据写入功能正常运行
send_email函数无任何输出,既不发送邮件也不抛出错误- 本地使用DirectRunner运行时所有功能正常,切换到DataflowRunner后该函数失效
完整代码
import apache_beam as beam import os import argparse import logging import pandas as pd import datetime import pytz from oauth2client.client import GoogleCredentials from datetime import datetime,date,timedelta from apache_beam.options.pipeline_options import PipelineOptions, StandardOptions import time import json import requests from google.cloud import pubsub_v1 import google.cloud.dlp_v2 from json import loads import smtplib from email.utils import formataddr from google.protobuf import json_format class readandwrite(beam.DoFn): def send_email(self,content_json): smtp_email = "gcpdev3@gmail.com" smtp_password = "password" email_subject_bigquery = "PII Information Detected!" email_body_bigquery = f''' Dear BigQuery Admin, We want to inform you that personally identifiable information (PII) has been detected in a transaction. - Duration: {content_json["duration"]} seconds - Text: "{content_json["text"]} Appropriate actions has been taken by the DLP teammsecure the data. Thank you, DataLake DLP Admin Team, sush Ltd. ''' try: with smtplib.SMTP("smtp.gmail.com", 587) as server: server.starttls() server.login(smtp_email, smtp_password) sender_name = "DLP DataLake Admin Team" sender_address = "no-reply@dlpadmin.com" formatted_sender = formataddr((sender_name, sender_address)) email_message_bigquery = f"From: {formatted_sender}\nSubject: {email_subject_bigquery}\nTo: sandeep.mohanty1998@gmail.com\n\n{email_body_bigquery}" server.sendmail(smtp_email, "sandeep@gmail.com", email_message_bigquery) print("Emails sent successfully.") except Exception as e: print("Error sending emails:", str(e)) def deidentify_content_with_dlp(self,content_json): import google.cloud.dlp_v2 dlp_client = google.cloud.dlp_v2.DlpServiceClient() item = {"value": json.dumps(content_json)} credit_card_info_type = {"name": "CREDIT_CARD_NUMBER"} phone_number_info_type = {"name": "PHONE_NUMBER"} deidentify_config = { "info_type_transformations": { "transformations": [ { "info_types": [credit_card_info_type], "primitive_transformation": { "character_mask_config": { "masking_character": "#", "number_to_mask": 7, "reverse_order": True, } }, }, { "info_types": [phone_number_info_type], "primitive_transformation": { "character_mask_config": { "masking_character": "#", "number_to_mask": 7, "reverse_order": True, } } }, ] } } project_id = "sandeepdev" parent = f"projects/{project_id}" try: response = dlp_client.deidentify_content( request={ "parent": parent, "deidentify_config": deidentify_config, "item": item, } ) logging.info("Applying DLP de-identification...") if response.item.value and response.overview.transformation_summaries: print("PII information found. Creating ServiceNow incident.") self.send_email(content_json) print("Email send to BigQuery Admin") return json.loads(response.item.value) if response.item.value else content_json except Exception as e: print(f"Error: {e}") print("Error during de-identification. Inserting original content to BigQuery.") return content_json logging.info("DLP de-identification completed...") def process(self, conetxt): import time import json import requests from google.cloud import pubsub_v1 import google.cloud.dlp_v2 from json import loads import smtplib from email.utils import formataddr from google.cloud import bigquery project_id = "sandeepdev" subscription_id = "audio_msg-sub" client_bigquery = bigquery.Client() subscriber = pubsub_v1.SubscriberClient() subscription_path = subscriber.subscription_path(project_id, subscription_id) masked_table_id = "sandeepdev.call_interaction.cleansed_dlp_raw_table" raw_table_id="sandeepdev.call_interaction.raw_audio_data" dlp_count_table_id="sandeepdev.call_interaction.dlp_count" max_messages = 1 count=0 rows_to_insert_raw = [] rows_to_insert_masked = [] rows_to_insert_dlp_count=[] logging.info("Starting data processing workflow...") while True: response = subscriber.pull(request={"subscription": subscription_path, "max_messages": max_messages}) for received_message in response.received_messages: message = received_message.message content_json = json.loads(message.data.decode('utf-8')) masked_data = self.deidentify_content_with_dlp(content_json) logging.info(masked_data) audio_file_name=masked_data['audio_file_name'] print("masked data is :" , masked_data) dlp_list= masked_data['text'].split(" ") for value in dlp_list: if '#' in value: count += 1 print("masked data count:", count) insert_data = { "audio_file_name": audio_file_name, "PII_count": count } rows_to_insert_dlp_count.append(insert_data) load_PII_count = client_bigquery.insert_rows_json(dlp_count_table_id, rows_to_insert_dlp_count) count=0 subscriber.acknowledge(request={"subscription": subscription_path, "ack_ids": [received_message.ack_id]}) rows_to_insert_raw.append(content_json) load_raw = client_bigquery.insert_rows_json(raw_table_id, rows_to_insert_raw) rows_to_insert_masked.append(masked_data) load_masked = client_bigquery.insert_rows_json(masked_table_id, rows_to_insert_masked) rows_to_insert_raw = [] rows_to_insert_masked = [] rows_to_insert_dlp_count=[] logging.info("Data processing workflow completed.") time.sleep(2) def run(): try: parser = argparse.ArgumentParser() parser.add_argument( '--dfBucket', required=True, help= ('Bucket where JARS/JDK is present') ) known_args, pipeline_args = parser.parse_known_args() global df_Bucket df_Bucket = known_args.dfBucket logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s') pipeline_options = PipelineOptions(pipeline_args) pipeline_options.view_as(StandardOptions).streaming = True pcoll = beam.Pipeline(options=pipeline_options) logging.info("Pipeline Starts") dummy= pcoll | 'Initializing..' >> beam.Create(['1']) dummy_env = dummy | 'Processing' >> beam.ParDo(readandwrite()) p=pcoll.run() logging.info('Job Run Successfully!') p.wait_until_finish() except Exception as e: logging.exception('Failed to launch datapipeline') logging.exception('Failed to launch datapipeline: %s', str(e)) raise if __name__ == '__main__': run()
运行命令
python3 /home/sandeepdev751/dlp_dataflow.py --runner DataflowRunner --project sandeepdev --region us-central1 --job_name dlpcheckerv5 --temp_location gs://sandeepdev/temp --dfBucket "sandeepdev" --setup_file /home/sandeepdev751/setup.py
解决方案
1. 替换print为logging,排查具体错误
DataFlow Worker中的print输出不会同步到控制台日志,改用logging模块可以在DataFlow的日志页面查看错误详情:
修改send_email函数:
def send_email(self,content_json): smtp_email = "gcpdev3@gmail.com" smtp_password = "password" email_subject_bigquery = "PII Information Detected!" email_body_bigquery = f''' Dear BigQuery Admin, We want to inform you that personally identifiable information (PII) has been detected in a transaction. - Duration: {content_json["duration"]} seconds - Text: "{content_json["text"]} Appropriate actions has been taken by the DLP team to secure the data. Thank you, DataLake DLP Admin Team, sush Ltd. ''' try: with smtplib.SMTP("smtp.gmail.com", 587) as server: server.starttls() server.login(smtp_email, smtp_password) sender_name = "DLP DataLake Admin Team" sender_address = "no-reply@dlpadmin.com" formatted_sender = formataddr((sender_name, sender_address)) email_message_bigquery = f"From: {formatted_sender}\nSubject: {email_subject_bigquery}\nTo: sandeep.mohanty1998@gmail.com\n\n{email_body_bigquery}" server.sendmail(smtp_email, "sandeep@gmail.com", email_message_bigquery) logging.info("Emails sent successfully.") except Exception as e: logging.error("Error sending emails:", exc_info=True)
2. 解决Gmail登录验证问题
- 如果开启了两步验证,需要创建应用专用密码替代普通密码
- 如果未开启两步验证,暂时允许低安全性应用访问(生产环境不推荐,建议改用GCP Cloud Functions或第三方邮件服务发送邮件)
3. 配置DataFlow Worker网络访问
确保DataFlow Worker所在的VPC防火墙规则允许出站流量到smtp.gmail.com的587端口:
- 目标IP:
smtp.gmail.com对应的IP段 - 协议:TCP,端口587
4. 优化Pipeline结构(最佳实践)
当前代码在DoFn中手动拉取Pub/Sub消息,不符合DataFlow的设计模式,改用Beam原生的Pub/Sub IO和BigQuery IO,提升可靠性和可维护性:
def run(): try: parser = argparse.ArgumentParser() parser.add_argument( '--dfBucket', required=True, help= ('Bucket where dependencies are stored') ) parser.add_argument( '--subscription', default="projects/sandeepdev/subscriptions/audio_msg-sub", help= 'Pub/Sub subscription full path' ) known_args, pipeline_args = parser.parse_known_args() logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s') pipeline_options = PipelineOptions(pipeline_args) pipeline_options.view_as(StandardOptions).streaming = True pcoll = beam.Pipeline(options=pipeline_options) logging.info("Pipeline Starts") # 读取Pub/Sub消息 messages = pcoll | 'Read Pub/Sub Messages' >> beam.io.ReadFromPubSub(subscription=known_args.subscription) parsed_messages = messages | 'Parse JSON' >> beam.Map(json.loads) # DLP脱敏处理,同时触发邮件 processed_messages = parsed_messages | 'DLP De-identify & Alert' >> beam.ParDo(readandwrite()) # 写入原始数据到BigQuery parsed_messages | 'Write Raw Data' >> beam.io.WriteToBigQuery( table="sandeepdev.call_interaction.raw_audio_data", write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED ) # 写入脱敏数据到BigQuery processed_messages | 'Write Masked Data' >> beam.io.WriteToBigQuery( table="sandeepdev.call_interaction.cleansed_dlp_raw_table", write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED ) p=pcoll.run() logging.info('Job Run Successfully!') p.wait_until_finish() except Exception as e: logging.exception('Failed to launch datapipeline') raise
内容的提问来源于stack exchange,提问作者Sandeep Mohanty
相关产品推荐
相关产品推荐

