You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.13 16:37:03