Zelfstudie: Operator voor e-mailzender van Gmail

In deze zelfstudie maakt u een python-run-function operator voor Lakeflow Designer waarmee de inhoud van een DataFrame als CSV-bijlage wordt verzonden via Gmail. In dit voorbeeld leert u hoe u op YAML gebaseerde operators bouwt die bijwerkingen uitvoeren, zoals het verzenden van meldingen of schrijven naar externe systemen. Zie Door de gebruiker gedefinieerde operators in Lakeflow Designer voor meer informatie.

Requirements

  • Een Azure Databricks-werkruimte met toegang om secret scopes te maken.
  • Een Gmail-account met een Google App-wachtwoord (vereist wanneer meervoudige verificatie (MFA) is ingeschakeld).
  • De Databricks CLI geïnstalleerd op uw lokale ontwikkelcomputer.

Stap 1: Geheime gegevens instellen

Sla uw Gmail-referenties op in een Azure Databricks geheim bereik, zodat de operator deze tijdens runtime kan ophalen.

  1. Maak een geheim bereik met behulp van de Azure Databricks CLI:

    databricks secrets create-scope my_email_scope
    
  2. Sla uw Gmail-appwachtwoord op in de scope:

    databricks secrets put-secret my_email_scope gmail_app_password
    

    U wordt gevraagd de geheime waarde in te voeren. Plak uw Gmail-appwachtwoord en sla het op.

Stap 2: De run() functie schrijven

Voor het python-run-function operatortype is een run() functie met deze handtekening vereist:

def run(config: Dict[str, Any], inputs: Dict[str, Any], spark) -> Dict[str, Any]:
  • config: Configuratiewaarden die door de gebruiker worden verstrekt in de gebruikersinterface van Lakeflow Designer.
  • inputs: Invoer-DataFrames met poortnaam als sleutel.
  • spark: De actieve Spark-sessie.

De functie moet een woordenlijst met uitvoerdataframes retourneren die zijn gekoppeld aan de naam van de uitvoerpoort.

Definieer en test de functie in een notebookcel:

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}

Stap 3: De functie testen

Test de functie met een voorbeeld van een 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

De secret_scope waarden secret_key in de configuratie zijn de namen van het geheime bereik en de sleutel die u in stap 1 hebt gemaakt, niet het werkelijke wachtwoord. De operator gebruikt deze namen om tijdens de uitvoering het wachtwoord op te halen uit de secrets van Azure Databricks.

Important

Test eerst met is_preview ingesteld op True om het pass-throughgedrag te controleren zonder e-mails te verzenden. Wanneer u klaar bent om de werkelijke e-mail te testen, stelt u deze in op is_previewFalse.

Stap 4: de YAML-definitie bouwen

Maak een bestand gmail_email_sender.yaml met de naam met de volgende inhoud:

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}

Stap 5: de operator opslaan en registreren

  1. Sla het YAML-bestand op in uw Azure Databricks werkruimte. Voorbeeld:

    /Workspace/Users/<user-name>/gmail_email_sender.yaml
    
  2. Voeg de operator toe aan uw .user_defined_operators.yaml bestand:

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

Zie Uw operator detecteerbaar maken voor meer informatie over registratieopties.

Permissions

Gebruikers die een werkstroom uitvoeren die deze operator bevat, moeten toegang hebben READ tot het geheime bereik of ze kunnen hun eigen geheime bereik en sleutelwaarden opgeven in de operatorconfiguratie. Gebruikers hebben ook leestoegang nodig tot het YAML-bestand in de werkruimte.

Toegang verlenen tot het geheime bereik:

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

Operator gebruiken in Lakeflow Designer

Na de registratie wordt de operator weergegeven in Lakeflow Designer met een invoerpoort voor uw gegevensbron en configuratievelden voor e-mail van afzender, geheim bereik, geheime sleutel, geadresseerden, onderwerp en hoofdtekst.

Wanneer de werkstroom wordt uitgevoerd, converteert de operator het dataframe naar CSV, voegt deze toe aan een e-mailbericht en verzendt deze naar elke geadresseerde. Het DataFrame wordt ongewijzigd doorgegeven aan de uitvoerpoort, zodat u aanvullende operators downstream kunt koppelen. Tijdens de voorbeeldweergave van de werkstroom wordt er geen e-mail verzonden.