Självstudie: Gmail-e-postsändaroperator

I den här självstudien skapar du en python-run-function operatör för Lakeflow Designer som skickar innehållet i en DataFrame som en CSV-bifogad fil via Gmail. Använd det här exemplet om du vill lära dig hur du skapar YAML-baserade operatorer som utför biverkningar, till exempel att skicka meddelanden eller skriva till externa system. Mer information finns i Användardefinierade operatorer i Lakeflow Designer.

Krav

  • En Azure Databricks arbetsyta med åtkomst för att skapa hemliga omfång.
  • Ett Gmail-konto med ett Google App-lösenord (krävs när multifaktorautentisering (MFA) är aktiverat).
  • Databricks CLI installerat på din lokala utvecklingsdator.

Steg 1: Konfigurera hemligheter

Lagra dina Gmail-autentiseringsuppgifter i ett hemlighetsomfång i Azure Databricks så att operatören kan hämta dem under körning.

  1. Skapa ett hemligt omfång med hjälp av Azure Databricks CLI:

    databricks secrets create-scope my_email_scope
    
  2. Lagra Gmail-appens lösenord inom omfattningen:

    databricks secrets put-secret my_email_scope gmail_app_password
    

    Du uppmanas att ange det hemliga värdet. Klistra in lösenordet för Gmail-appen och spara.

Steg 2: Skriv run() funktionen

Operatortypen python-run-function kräver en run() funktion med den här signaturen:

def run(config: Dict[str, Any], inputs: Dict[str, Any], spark) -> Dict[str, Any]:
  • config: Konfigurationsvärden som tillhandahålls av användaren i Lakeflow Designer-användargränssnittet.
  • inputs: Indataramar som är nyckelade efter portnamn.
  • spark: Den aktiva Spark-sessionen.

Funktionen måste returnera en ordlista med utdataramar som är nyckelade efter portnamnet för utdata.

Definiera och testa funktionen i en notebook-cell:

from typing import Dict, Any

def run(config: Dict[str, Any], inputs: Dict[str, Any], spark) -> Dict[str, Any]:
    input_df = inputs["data"]

    # Skip side effects during Designer preview
    if config.get("is_preview", False):
        return {"data": input_df}

    import smtplib
    import os
    from email.mime.multipart import MIMEMultipart
    from email.mime.text import MIMEText
    from email.mime.base import MIMEBase
    from email import encoders

    sender_email = config.get("sender_email", "")
    secret_scope = config.get("secret_scope", "")
    secret_key = config.get("secret_key", "")
    recipients_raw = config.get("recipients", "")
    subject = config.get("subject", "")
    body = config.get("body", "")

    if not sender_email:
        raise ValueError("Sender Email is required.")
    if not secret_scope or not secret_key:
        raise ValueError("Secret Scope and Secret Key are required.")
    if not recipients_raw:
        raise ValueError("At least one recipient is required.")

    recipients = [r.strip() for r in recipients_raw.split(",") if r.strip()]
    if not recipients:
        raise ValueError("At least one valid recipient email is required.")

    # Retrieve password from Databricks secrets
    from pyspark.dbutils import DBUtils
    dbutils = DBUtils(spark)
    sender_password = dbutils.secrets.get(scope=secret_scope, key=secret_key)

    # Convert DataFrame to CSV
    pdf = input_df.toPandas()
    file_path = "/tmp/designer_email_attachment.csv"
    pdf.to_csv(file_path, index=False)

    # Send email to each recipient
    for recipient in recipients:
        msg = MIMEMultipart()
        msg["From"] = sender_email
        msg["To"] = recipient
        msg["Subject"] = subject
        msg.attach(MIMEText(body, "plain"))

        with open(file_path, "rb") as attachment:
            part = MIMEBase("application", "octet-stream")
            part.set_payload(attachment.read())
            encoders.encode_base64(part)
            part.add_header(
                "Content-Disposition",
                f"attachment; filename={os.path.basename(file_path)}",
            )
            msg.attach(part)

        with smtplib.SMTP_SSL("smtp.gmail.com", 465) as server:
            server.login(sender_email, sender_password)
            server.send_message(msg)

    # Clean up temp file
    if os.path.exists(file_path):
        os.remove(file_path)

    return {"data": input_df}

Steg 3: Testa funktionen

Testa funktionen med ett exempel på DataFrame:

test_df = spark.createDataFrame(
    [("Alice", 100), ("Bob", 200)],
    ["name", "amount"]
)

# Test in preview mode (no email sent)
result = run(
    config={
        "is_preview": True,
        "sender_email": "you@gmail.com",
        "secret_scope": "my_email_scope",
        "secret_key": "gmail_app_password",
        "recipients": "alice@example.com",
        "subject": "Test",
        "body": "Test body"
    },
    inputs={"data": test_df},
    spark=spark
)

result["data"].show()
# Expected: the original DataFrame, unchanged

Note

Värdena secret_scope och secret_key i konfigurationen är namnen på det hemliga omfånget och nyckeln som du skapade i steg 1 – inte det faktiska lösenordet. Operatorn använder dessa namn för att hämta lösenordet från hemliga nycklar i Azure Databricks under körning.

Important

Testa först med is_preview inställt på True för att verifiera vidarebefordringsbeteendet utan att skicka någon e-post. När du är redo att testa det faktiska e-postmeddelandet anger du is_preview till False.

Steg 4: Skapa YAML-definitionen

Skapa en fil med gmail_email_sender.yaml namnet med följande innehåll:

schema: user-defined-operator-v0.1.0
id: gmail_email_sender
type: python-run-function
version: '1.0.0'
name: Gmail Email Sender
description: Sends the input DataFrame as a CSV attachment via Gmail SMTP to one or more recipients.

config:
  type: object
  properties:
    is_preview:
      type: boolean
      format: is_preview
      default: false
    sender_email:
      type: string
      title: Sender Email
      default: ''
      examples:
        - 'you@gmail.com'
      x-ui:
        widget: input
    secret_scope:
      type: string
      title: Secret Scope
      default: ''
      examples:
        - 'my_email_scope'
      x-ui:
        widget: input
    secret_key:
      type: string
      title: Secret Key
      default: ''
      examples:
        - 'gmail_app_password'
      x-ui:
        widget: input
    recipients:
      type: string
      title: Recipients
      default: ''
      examples:
        - 'alice@example.com, bob@example.com'
      x-ui:
        widget: textarea
        rows: 2
    subject:
      type: string
      title: Subject
      default: ''
      examples:
        - 'Designer Output Data'
      x-ui:
        widget: input
    body:
      type: string
      title: Email Body
      default: "Hello,\n\nAttached is the latest data.\n\nBest,\nDatabricks Workflow"
      x-ui:
        widget: textarea
        rows: 6
  required:
    - sender_email
    - secret_scope
    - secret_key
    - recipients
    - subject
  additionalProperties: false

ports:
  input:
    - name: data
      title: Input Data
      mime: application/vnd.databricks.dataframe
  output:
    - name: data
      title: Output Data
      mime: application/vnd.databricks.dataframe

run_function:
  type: inline
  code: |
    from typing import Dict, Any

    def run(config: Dict[str, Any], inputs: Dict[str, Any], spark) -> Dict[str, Any]:
        input_df = inputs["data"]

        if config.get("is_preview", False):
            return {"data": input_df}

        import smtplib
        import os
        from email.mime.multipart import MIMEMultipart
        from email.mime.text import MIMEText
        from email.mime.base import MIMEBase
        from email import encoders

        sender_email = config.get("sender_email", "")
        secret_scope = config.get("secret_scope", "")
        secret_key = config.get("secret_key", "")
        recipients_raw = config.get("recipients", "")
        subject = config.get("subject", "")
        body = config.get("body", "")

        if not sender_email:
            raise ValueError("Sender Email is required.")
        if not secret_scope or not secret_key:
            raise ValueError("Secret Scope and Secret Key are required.")
        if not recipients_raw:
            raise ValueError("At least one recipient is required.")

        recipients = [r.strip() for r in recipients_raw.split(",") if r.strip()]
        if not recipients:
            raise ValueError("At least one valid recipient email is required.")

        from pyspark.dbutils import DBUtils
        dbutils = DBUtils(spark)
        sender_password = dbutils.secrets.get(scope=secret_scope, key=secret_key)

        pdf = input_df.toPandas()
        file_path = "/tmp/designer_email_attachment.csv"
        pdf.to_csv(file_path, index=False)

        for recipient in recipients:
            msg = MIMEMultipart()
            msg["From"] = sender_email
            msg["To"] = recipient
            msg["Subject"] = subject
            msg.attach(MIMEText(body, "plain"))

            with open(file_path, "rb") as attachment:
                part = MIMEBase("application", "octet-stream")
                part.set_payload(attachment.read())
                encoders.encode_base64(part)
                part.add_header(
                    "Content-Disposition",
                    f"attachment; filename={os.path.basename(file_path)}",
                )
                msg.attach(part)

            with smtplib.SMTP_SSL("smtp.gmail.com", 465) as server:
                server.login(sender_email, sender_password)
                server.send_message(msg)

        if os.path.exists(file_path):
            os.remove(file_path)

        return {"data": input_df}

Steg 5: Spara och registrera operatorn

  1. Spara YAML-filen på din Azure Databricks arbetsyta. Ett exempel:

    /Workspace/Users/<user-name>/gmail_email_sender.yaml
    
  2. Lägg till operatorn i din .user_defined_operators.yaml-fil:

    operators:
      - /Workspace/Users/<user-name>/gmail_email_sender.yaml
    

Mer information om registreringsalternativ finns i Gör operatören identifierbar.

Permissions

Användare som kör ett arbetsflöde som innehåller den här operatorn behöver READ åtkomst till det hemliga omfånget, eller så kan de ange sitt eget hemliga omfång och nyckelvärden i operatorkonfigurationen. Användarna behöver också läsbehörighet till YAML-filen på arbetsytan.

Så här beviljar du åtkomst till det hemliga omfånget:

databricks secrets put-acl my_email_scope <user-or-group> READ

Använda operatorn i Lakeflow Designer

Efter registreringen visas operatorn i Lakeflow Designer med en indataport för datakällan och konfigurationsfälten för avsändarens e-post, hemlighetsomfång, hemlig nyckel, mottagare, ämne och brödtext.

När arbetsflödet körs konverterar operatorn dataramen till CSV, kopplar den till ett e-postmeddelande och skickar den till varje mottagare. DataFramen skickas vidare oförändrad till utdataporten, så att du kan koppla på fler operatorer nedströms. Under förhandsversionen av arbetsflödet skickas inget e-postmeddelande.