Azure Functions 的 Apache Kafka 觸發程序

在 Azure Functions 中使用 Apache Kafka 觸發器,根據 Kafka 主題中的訊息執行你的函式程式碼。 您也可以使用 Kafka 輸出系結 ,從函式寫入主題。 如需安裝和組態詳細數據的詳細資訊,請參閱 Azure Functions 的 Apache Kafka 系結概觀。

重要

Kafka 綁定可用於 Flex Consumption 計畫、 Elastic Premium 計畫及 專用(App Service)方案的功能。 它們僅支援於 Functions 執行環境的 4.x 版本。

範例

目前這款綁定沒有 Go 支援。

觸發程式的使用方式取決於函式應用程式中所使用的 C# 形式,這可以是下列其中一種模式:

一個編譯好的 C# 函式,使用一個獨立的工作 者程序類別函式庫 ,執行在與執行時分開的程序中。

您使用的屬性取決於特定事件提供者。

下列範例示範 C# 函式,以 Kafka 事件的形式讀取和記錄 Kafka 訊息:

[Function("KafkaTrigger")]
public static void Run(
    [KafkaTrigger("BrokerList",
                  "topic",
                  Username = "ConfluentCloudUserName",
                  Password = "ConfluentCloudPassword",
                  Protocol = BrokerProtocol.SaslSsl,
                  AuthenticationMode = BrokerAuthenticationMode.Plain,
                  ConsumerGroup = "$Default")] string eventData, FunctionContext context)
{
    var logger = context.GetLogger("KafkaFunction");
    logger.LogInformation($"C# Kafka trigger function processed a message: {JObject.Parse(eventData)["Value"]}");
}

若要接收批次中的事件,請使用字串陣列做為輸入,如下列範例所示:

[Function("KafkaTriggerMany")]
public static void Run(
    [KafkaTrigger("BrokerList",
                  "topic",
                  Username = "ConfluentCloudUserName",
                  Password = "ConfluentCloudPassword",
                  Protocol = BrokerProtocol.SaslSsl,
                  AuthenticationMode = BrokerAuthenticationMode.Plain,
                  ConsumerGroup = "$Default",
                  IsBatched = true)] string[] events, FunctionContext context)
{
    foreach (var kevent in events)
    {
        var logger = context.GetLogger("KafkaFunction");
        logger.LogInformation($"C# Kafka trigger function processed a message: {JObject.Parse(kevent)["Value"]}");
    }

下列函式會記錄 Kafka 事件的訊息和標頭:

[Function("KafkaTriggerWithHeaders")]
public static void Run(
    [KafkaTrigger("BrokerList",
                  "topic",
                  Username = "ConfluentCloudUserName",
                  Password = "ConfluentCloudPassword",
                  Protocol = BrokerProtocol.SaslSsl,
                  AuthenticationMode = BrokerAuthenticationMode.Plain,
                  ConsumerGroup = "$Default")] string eventData, FunctionContext context)
{
    var eventJsonObject = JObject.Parse(eventData);
    var logger = context.GetLogger("KafkaFunction");
    logger.LogInformation($"C# Kafka trigger function processed a message: {eventJsonObject["Value"]}");
    var headersJArr = eventJsonObject["Headers"] as JArray;
    logger.LogInformation("Headers for this event: ");
    foreach (JObject header in headersJArr)
    {
        logger.LogInformation($"{header["Key"]} {System.Text.Encoding.UTF8.GetString((byte[])header["Value"])}");

    }
}

如需一組完整的工作 .NET 範例,請參閱 Kafka擴充功能存放庫。

觸發器的使用取決於你版本的 Node.js 程式設計模型。

在 Node.js v4 模型中,你直接在函式程式碼中定義觸發器。 如需詳細資訊,請參閱 Azure Functions Node.js 開發人員指南。

在這些例子中,事件提供者是 Confluent 或 Azure 事件中樞。 這些範例展示了如何定義讀取卡夫卡訊息函式的卡夫卡觸發器。

const { app } = require("@azure/functions");

async function kafkaTrigger(event, context) {
  context.log("Event Offset: " + event.Offset);
  context.log("Event Partition: " + event.Partition);
  context.log("Event Topic: " + event.Topic);
  context.log("Event Timestamp: " + event.Timestamp);
  context.log("Event Key: " + event.Key);
  context.log("Event Value (as string): " + event.Value);

  let event_obj = JSON.parse(event.Value);

  context.log("Event Value Object: ");
  context.log("   Value.registertime: ", event_obj.registertime.toString());
  context.log("   Value.userid: ", event_obj.userid);
  context.log("   Value.regionid: ", event_obj.regionid);
  context.log("   Value.gender: ", event_obj.gender);
}

app.generic("Kafkatrigger", {
  trigger: {
    type: "kafkaTrigger",
    direction: "in",
    name: "event",
    topic: "topic",
    brokerList: "%BrokerList%",
    username: "%ConfluentCloudUserName%",
    password: "%ConfluentCloudPassword%",
    consumerGroup: "$Default",
    protocol: "saslSsl",
    authenticationMode: "plain",
    dataType: "string"
  },
  handler: kafkaTrigger,
});

要接收批次事件,請將值設 cardinality 為 many,如以下範例所示:

const { app } = require("@azure/functions");

async function kafkaTriggerMany(events, context) {
  for (const event of events) {
    context.log("Event Offset: " + event.Offset);
    context.log("Event Partition: " + event.Partition);
    context.log("Event Topic: " + event.Topic);
    context.log("Event Key: " + event.Key);
    context.log("Event Timestamp: " + event.Timestamp);
    context.log("Event Value (as string): " + event.Value);

    let event_obj = JSON.parse(event.Value);

    context.log("Event Value Object: ");
    context.log("   Value.registertime: ", event_obj.registertime.toString());
    context.log("   Value.userid: ", event_obj.userid);
    context.log("   Value.regionid: ", event_obj.regionid);
    context.log("   Value.gender: ", event_obj.gender);
  }
}

app.generic("kafkaTriggerMany", {
  trigger: {
    type: "kafkaTrigger",
    direction: "in",
    name: "event",
    topic: "topic",
    brokerList: "%BrokerList%",
    username: "%ConfluentCloudUserName%",
    password: "%ConfluentCloudPassword%",
    consumerGroup: "$Default",
    protocol: "saslSsl",
    authenticationMode: "plain",
    dataType: "string",
    cardinality: "MANY"
  },
  handler: kafkaTriggerMany,
});

您可以定義傳遞至觸發程式之事件的一般 Avro 架構 。 此範例定義了特定提供者的觸發條件,並採用通用的 Avro 架構:

const { app } = require("@azure/functions");

async function kafkaAvroGenericTrigger(event, context) {
  context.log("Processed kafka event: ", event);
  if (context.triggerMetadata?.key !== undefined) {
    context.log("message key: ", context.triggerMetadata?.key);
  }
}

app.generic("kafkaAvroGenericTrigger", {
  trigger: {
    type: "kafkaTrigger",
    direction: "in",
    name: "event",
    protocol: "SASLSSL",
    password: "EventHubConnectionString",
    dataType: "string",
    topic: "topic",
    authenticationMode: "PLAIN",
    avroSchema:
      '{"type":"record","name":"Payment","namespace":"io.confluent.examples.clients.basicavro","fields":[{"name":"id","type":"string"},{"name":"amount","type":"double"},{"name":"type","type":"string"}]}',
    consumerGroup: "$Default",
    username: "$ConnectionString",
    brokerList: "%BrokerList%",
  },
  handler: kafkaAvroGenericTrigger,
});

如需一組完整的 JavaScript 範例,請參閱 Kafka 擴充功能存放庫。

import { app, InvocationContext } from "@azure/functions";

// This is a sample interface that describes the actual data in your event.
interface EventData {
  registertime: number;
  userid: string;
  regionid: string;
  gender: string;
}

export async function kafkaTrigger(
  event: any,
  context: InvocationContext
): Promise<void> {
  context.log("Event Offset: " + event.Offset);
  context.log("Event Partition: " + event.Partition);
  context.log("Event Topic: " + event.Topic);
  context.log("Event Timestamp: " + event.Timestamp);
  context.log("Event Value (as string): " + event.Value);

  let event_obj: EventData = JSON.parse(event.Value);

  context.log("Event Value Object: ");
  context.log("   Value.registertime: ", event_obj.registertime.toString());
  context.log("   Value.userid: ", event_obj.userid);
  context.log("   Value.regionid: ", event_obj.regionid);
  context.log("   Value.gender: ", event_obj.gender);
}

app.generic("Kafkatrigger", {
  trigger: {
    type: "kafkaTrigger",
    direction: "in",
    name: "event",
    topic: "topic",
    brokerList: "%BrokerList%",
    username: "%ConfluentCloudUserName%",
    password: "%ConfluentCloudPassword%",
    consumerGroup: "$Default",
    protocol: "saslSsl",
    authenticationMode: "plain",
    dataType: "string"
  },
  handler: kafkaTrigger,
});

要接收批次事件,請將值設 cardinality 為 many,如以下範例所示:

import { app, InvocationContext } from "@azure/functions";

// This is a sample interface that describes the actual data in your event.
interface EventData {
    registertime: number;
    userid: string;
    regionid: string;
    gender: string;
}

interface KafkaEvent {
    Offset: number;
    Partition: number;
    Topic: string;
    Timestamp: number;
    Value: string;
}

export async function kafkaTriggerMany(
    events: any,
    context: InvocationContext
): Promise<void> {
    for (const event of events) {
        context.log("Event Offset: " + event.Offset);
        context.log("Event Partition: " + event.Partition);
        context.log("Event Topic: " + event.Topic);
        context.log("Event Timestamp: " + event.Timestamp);
        context.log("Event Value (as string): " + event.Value);

        let event_obj: EventData = JSON.parse(event.Value);

        context.log("Event Value Object: ");
        context.log("   Value.registertime: ", event_obj.registertime.toString());
        context.log("   Value.userid: ", event_obj.userid);
        context.log("   Value.regionid: ", event_obj.regionid);
        context.log("   Value.gender: ", event_obj.gender);
    }
}

app.generic("kafkaTriggerMany", {
  trigger: {
    type: "kafkaTrigger",
    direction: "in",
    name: "event",
    topic: "topic",
    brokerList: "%BrokerList%",
    username: "%ConfluentCloudUserName%",
    password: "%ConfluentCloudPassword%",
    consumerGroup: "$Default",
    protocol: "saslSsl",
    authenticationMode: "plain",
    dataType: "string",
    cardinality: "MANY"
  },
  handler: kafkaTriggerMany,
});

您可以定義傳遞至觸發程式之事件的一般 Avro 架構 。 此範例定義了特定提供者的觸發條件,並採用通用的 Avro 架構:

import { app, InvocationContext } from "@azure/functions";

export async function kafkaAvroGenericTrigger(
  event: any,
  context: InvocationContext
): Promise<void> {
  context.log("Processed kafka event: ", event);
  context.log(
    `Message ID: ${event.id}, amount: ${event.amount}, type: ${event.type}`
  );
  if (context.triggerMetadata?.key !== undefined) {
    context.log(`Message Key : ${context.triggerMetadata?.key}`);
  }
}

app.generic("kafkaAvroGenericTrigger", {
  trigger: {
    type: "kafkaTrigger",
    direction: "in",
    name: "event",
    protocol: "SASLSSL",
    username: "ConfluentCloudUsername",
    password: "ConfluentCloudPassword",
    dataType: "string",
    topic: "topic",
    authenticationMode: "PLAIN",
    avroSchema:
      '{"type":"record","name":"Payment","namespace":"io.confluent.examples.clients.basicavro","fields":[{"name":"id","type":"string"},{"name":"amount","type":"double"},{"name":"type","type":"string"}]}',
    consumerGroup: "$Default",
    brokerList: "%BrokerList%",
  },
  handler: kafkaAvroGenericTrigger,
});

欲了解完整的 TypeScript 範例,請參閱 Kafka 擴充套件庫。

檔案的 function.json 具體屬性取決於你的活動提供者。 在這些例子中,事件提供者是 Confluent 或 Azure 事件中樞。 下列範例顯示讀取和記錄 Kafka 訊息之函式的 Kafka 觸發程式。

以下 function.json 檔案定義了特定提供者的觸發條件:

{
    "bindings": [
      {
            "type": "kafkaTrigger",
            "name": "kafkaEvent",
            "direction": "in",
            "protocol" : "SASLSSL",
            "password" : "%ConfluentCloudPassword%",
            "dataType" : "string",
            "topic" : "topic",
            "authenticationMode" : "PLAIN",
            "consumerGroup" : "$Default",
            "username" : "%ConfluentCloudUserName%",
            "brokerList" : "%BrokerList%",
            "sslCaLocation": "confluent_cloud_cacert.pem"
        }
    ]
}

當函式被觸發時,以下程式碼會執行:

using namespace System.Net

param($kafkaEvent, $TriggerMetadata)

Write-Output "Powershell Kafka trigger function called for message $kafkaEvent.Value"

若要在批次中接收事件,請將 cardinality function.json 檔案中的 值設定為 many ,如下列範例所示:

{
    "bindings": [
      {
            "type": "kafkaTrigger",
            "name": "kafkaEvent",
            "direction": "in",
            "protocol" : "SASLSSL",
            "password" : "%ConfluentCloudPassword%",
            "dataType" : "string",
            "topic" : "topic",
            "authenticationMode" : "PLAIN",
            "cardinality" : "MANY",
            "consumerGroup" : "$Default",
            "username" : "%ConfluentCloudUserName%",
            "brokerList" : "%BrokerList%",
            "sslCaLocation": "confluent_cloud_cacert.pem"
        }
    ]
}

以下程式碼解析事件陣列並記錄事件資料:

using namespace System.Net

param($kafkaEvents, $TriggerMetadata)

$kafkaEvents
foreach ($kafkaEvent in $kafkaEvents) {
    $event = $kafkaEvent | ConvertFrom-Json -AsHashtable
    Write-Output "Powershell Kafka trigger function called for message $event.Value"
}

以下程式碼記錄標頭資料:

using namespace System.Net

param($kafkaEvents, $TriggerMetadata)

foreach ($kafkaEvent in $kafkaEvents) {
    $kevent = $kafkaEvent | ConvertFrom-Json -AsHashtable
    Write-Output "Powershell Kafka trigger function called for message $kevent.Value"
    Write-Output "Headers for this message:"
    foreach ($header in $kevent.Headers) {
        $DecodedValue = [System.Text.Encoding]::Unicode.GetString([System.Convert]::FromBase64String($header.Value))
        $Key = $header.Key
        Write-Output "Key: $Key Value: $DecodedValue"
    }
}

您可以定義傳遞至觸發程式之事件的一般 Avro 架構 。 下列function.json會使用一般 Avro 架構來定義特定提供者的觸發程式:

{
  "bindings" : [ {
    "type" : "kafkaTrigger",
    "direction" : "in",
    "name" : "kafkaEvent",
    "protocol" : "SASLSSL",
    "password" : "ConfluentCloudPassword",
    "topic" : "topic",
    "authenticationMode" : "PLAIN",
    "avroSchema" : "{\"type\":\"record\",\"name\":\"Payment\",\"namespace\":\"io.confluent.examples.clients.basicavro\",\"fields\":[{\"name\":\"id\",\"type\":\"string\"},{\"name\":\"amount\",\"type\":\"double\"},{\"name\":\"type\",\"type\":\"string\"}]}",
    "consumerGroup" : "$Default",
    "username" : "ConfluentCloudUsername",
    "brokerList" : "%BrokerList%"
  } ]
}

當函式被觸發時,以下程式碼會執行:

using namespace System.Net

param($kafkaEvent, $TriggerMetadata)

Write-Output "Powershell Kafka trigger function called for message $kafkaEvent.Value"

如需一組完整的運作 PowerShell 範例,請參閱 Kafka 擴充功能存放庫。

觸發器的使用取決於你版本的 Python 程式設計模型。

在 Python v2 模型中,你直接用 decorator 在函式程式碼中定義觸發器。 更多資訊請參閱 Azure Functions Python 開發者指南。

這些範例展示了如何定義讀取卡夫卡訊息函式的卡夫卡觸發器。

@KafkaTrigger.function_name(name="KafkaTrigger")
@KafkaTrigger.kafka_trigger(
    arg_name="kevent",
    topic="KafkaTopic",
    broker_list="KafkaBrokerList",
    username="KafkaUsername",
    password="KafkaPassword",
    protocol="SaslSsl",
    authentication_mode="Plain",
    consumer_group="$Default1")
def kafka_trigger(kevent : func.KafkaEvent):
    logging.info(kevent.get_body().decode('utf-8'))
    logging.info(kevent.metadata)

此範例透過將值設定 cardinality 為 many來接收批次事件。

@KafkaTrigger.function_name(name="KafkaTriggerMany")
@KafkaTrigger.kafka_trigger(
    arg_name="kevents",
    topic="KafkaTopic",
    broker_list="KafkaBrokerList",
    username="KafkaUsername",
    password="KafkaPassword",
    protocol="SaslSsl",
    authentication_mode="Plain",
    cardinality="MANY",
    data_type="string",
    consumer_group="$Default2")
def kafka_trigger_many(kevents : typing.List[func.KafkaEvent]):
    for event in kevents:
        logging.info(event.get_body())

您可以定義傳遞至觸發程式之事件的一般 Avro 架構 。

@KafkaTriggerAvro.function_name(name="KafkaTriggerAvroOne")
@KafkaTriggerAvro.kafka_trigger(
    arg_name="kafkaTriggerAvroGeneric",
    topic="KafkaTopic",
    broker_list="KafkaBrokerList",
    username="KafkaUsername",
    password="KafkaPassword",
    protocol="SaslSsl",
    authentication_mode="Plain",
    consumer_group="$Default",
    avro_schema= "{\"type\":\"record\",\"name\":\"Payment\",\"namespace\":\"io.confluent.examples.clients.basicavro\",\"fields\":[{\"name\":\"id\",\"type\":\"string\"},{\"name\":\"amount\",\"type\":\"double\"},{\"name\":\"type\",\"type\":\"string\"}]}")
def kafka_trigger_avro_one(kafkaTriggerAvroGeneric : func.KafkaEvent):
    logging.info(kafkaTriggerAvroGeneric.get_body().decode('utf-8'))
    logging.info(kafkaTriggerAvroGeneric.metadata)

如需一組完整的工作 Python 範例,請參閱 Kafka 擴充功能存放庫。

您用來設定觸發程式的批註取決於特定事件提供者。

下列範例顯示 Java 函式,可讀取和記錄 Kafka 事件的內容:

@FunctionName("KafkaTrigger")
public void runSingle(
        @KafkaTrigger(
            name = "KafkaTrigger",
            topic = "topic",  
            brokerList="%BrokerList%",
            consumerGroup="$Default", 
            username = "%ConfluentCloudUsername%", 
            password = "ConfluentCloudPassword",
            authenticationMode = BrokerAuthenticationMode.PLAIN,
            protocol = BrokerProtocol.SASLSSL,
            // sslCaLocation = "confluent_cloud_cacert.pem", // Enable this line for windows.
            dataType = "string"
         ) String kafkaEventData,
        final ExecutionContext context) {
        context.getLogger().info(kafkaEventData);
}

若要接收批次中的事件,請使用輸入字串作為陣列,如下列範例所示:

@FunctionName("KafkaTriggerMany")
public void runMany(
        @KafkaTrigger(
            name = "kafkaTriggerMany",
            topic = "topic",  
            brokerList="%BrokerList%",
            consumerGroup="$Default", 
            username = "%ConfluentCloudUsername%", 
            password = "ConfluentCloudPassword",
            authenticationMode = BrokerAuthenticationMode.PLAIN,
            protocol = BrokerProtocol.SASLSSL,
            // sslCaLocation = "confluent_cloud_cacert.pem", // Enable this line for windows.
            cardinality = Cardinality.MANY,
            dataType = "string"
         ) String[] kafkaEvents,
        final ExecutionContext context) {
        for (String kevent: kafkaEvents) {
            context.getLogger().info(kevent);
        }    
}

下列函式會記錄 Kafka 事件的訊息和標頭:

@FunctionName("KafkaTriggerManyWithHeaders")
public void runSingle(
        @KafkaTrigger(
            name = "KafkaTrigger",
            topic = "topic",  
            brokerList="%BrokerList%",
            consumerGroup="$Default", 
            username = "%ConfluentCloudUsername%", 
            password = "ConfluentCloudPassword",
            authenticationMode = BrokerAuthenticationMode.PLAIN,
            protocol = BrokerProtocol.SASLSSL,
            // sslCaLocation = "confluent_cloud_cacert.pem", // Enable this line for windows.
            dataType = "string",
            cardinality = Cardinality.MANY
         ) List<String> kafkaEvents,
        final ExecutionContext context) {
            Gson gson = new Gson(); 
            for (String keventstr: kafkaEvents) {
                KafkaEntity kevent = gson.fromJson(keventstr, KafkaEntity.class);
                context.getLogger().info("Java Kafka trigger function called for message: " + kevent.Value);
                context.getLogger().info("Headers for the message:");
                for (KafkaHeaders header : kevent.Headers) {
                    String decodedValue = new String(Base64.getDecoder().decode(header.Value));
                    context.getLogger().info("Key:" + header.Key + " Value:" + decodedValue);                    
                }                
            }
        }

您可以定義傳遞至觸發程式之事件的一般 Avro 架構 。 下列函式會使用一般 Avro 架構來定義特定提供者的觸發程式:

private static final String schema = "{\"type\":\"record\",\"name\":\"Payment\",\"namespace\":\"io.confluent.examples.clients.basicavro\",\"fields\":[{\"name\":\"id\",\"type\":\"string\"},{\"name\":\"amount\",\"type\":\"double\"},{\"name\":\"type\",\"type\":\"string\"}]}";

@FunctionName("KafkaAvroGenericTrigger")
public void runOne(
        @KafkaTrigger(
                name = "kafkaAvroGenericSingle",
                topic = "topic",
                brokerList="%BrokerList%",
                consumerGroup="$Default",
                username = "ConfluentCloudUsername",
                password = "ConfluentCloudPassword",
                avroSchema = schema,
                authenticationMode = BrokerAuthenticationMode.PLAIN,
                protocol = BrokerProtocol.SASLSSL) Payment payment,
        final ExecutionContext context) {
    context.getLogger().info(payment.toString());
}

如需 Confluent 的完整工作 Java 範例集,請參閱 Kafka 擴充功能存放庫。

屬性

進程內和隔離的背景工作進程 C# 連結庫都會使用 KafkaTriggerAttribute 來定義函式觸發程式。

下表說明使用此觸發屬性可設定的屬性:

參數 描述
BrokerList (必要)觸發程式所監視的 Kafka 訊息代理程式清單。 如需詳細資訊,請參閱 連線 。
主題 (必要)觸發程式所監視的主題。
ConsumerGroup (選擇性)觸發程式所使用的 Kafka 取用者群組。
AvroSchema (可選)使用 Avro 協定時,訊息值的通用記錄結構。
KeyAvroSchema (可選)使用 Avro 協定時,訊息金鑰的通用記錄架構。
KeyDataType (可選)資料型別以接收訊息金鑰,如同卡夫卡主題。 若 KeyAvroSchema 設定為 ,則此值為通用記錄。 公認的值為 Int、 Long、 String和 Binary。
AuthenticationMode (選擇性)使用簡單驗證和安全性層 (SASL) 驗證時的驗證模式。 支援的值為NotSet(預設值)、Gssapi、 PlainScramSha256ScramSha512OAuthBearer和 。
使用者名稱 (選擇性)SASL 驗證的用戶名稱。 當 為AuthenticationMode時Gssapi不受支援。 如需詳細資訊,請參閱 連線 。
密碼 (選擇性)SASL 驗證的密碼。 當 為AuthenticationMode時Gssapi不受支援。 如需詳細資訊,請參閱 連線 。
通訊協定 (選擇性)與訊息代理程式通訊時所使用的安全性通訊協定。 支援的值為NotSet(預設值)、plaintext、sslsasl_plaintextsasl_ssl。
SslCaLocation (選擇性)用於驗證訊息代理程序憑證的 CA 憑證檔案路徑。
SslCertificateLocation (選擇性)用戶端憑證的路徑。
SslKeyLocation (選擇性)用於驗證之用戶端私鑰 (PEM) 的路徑。
SslKeyPassword (選擇性)用戶端憑證的密碼。
SslCertificatePEM (可選)以 PEM 格式呈現的用戶端憑證,作為字串。 如需詳細資訊,請參閱 連線 。
SslKeyPEM (可選)以 PEM 格式呈現的客戶端私鑰,作為字串。 如需詳細資訊,請參閱 連線 。
SslCaPEM (可選)CA 憑證以 PEM 格式以字串形式呈現。 如需詳細資訊,請參閱 連線 。
SslCertificate與KeyPEM (可選)客戶端憑證與金鑰以 PEM 格式呈現字串。 如需詳細資訊,請參閱 連線 。
SchemaRegistryUrl (可選)Avro 架構登錄檔的網址。 如需詳細資訊,請參閱 連線 。
SchemaRegistryUsername (可選)Avro 架構登錄檔的用戶名。 如需詳細資訊,請參閱 連線 。
SchemaRegistryPassword (可選)Avro 架構登錄檔的密碼。 如需詳細資訊,請參閱 連線 。
OAuthBearerMethod (可選)OAuth Bearer 方法。 接受的值是 oidc 和 default。
OAuthBearerClientId (可選)當 OAuthBearerMethod 設為 oidc時,這表示 OAuth 承載的客戶端 ID。 如需詳細資訊,請參閱 連線 。
OAuthBearerClientSecret (可選)當 OAuthBearerMethod 設定為 oidc時,這會指定 OAuth 承載的客戶端秘密。 如需詳細資訊,請參閱 連線 。
OAuthBearerScope(開放承載者望遠鏡) (可選)指定向經紀人申請存取的範圍。
OAuthBearerTokenEndpointUrl (可選)OAuth/OIDC 發行者令牌端點 HTTP(S) URI 用於在使用方法時 oidc 取得令牌。 如需詳細資訊,請參閱 連線 。
HttpsCaLocation (可選)用於驗證 OAuth/OIDC 令牌端點憑證的檔案或目錄路徑,指向 CA 憑證。 特殊值 probe 使用作業系統的預設憑證路徑。 僅由孤立工人模型支持。
HttpsCaPem (可選)CA 憑證用於驗證 PEM 格式的 OAuth/OIDC 令牌端點憑證。 僅由孤立工人模型支持。
OAuthBearer擴充 (可選)使用方法時 oidc ,提供逗號分隔的鍵=值對清單作為額外資訊,供經紀人使用。 例如: supportFeatureX=true,organizationId=sales-emea 。

重要

HttpsCaLocation HttpsCaPem目前動態比例計畫中還不支援這些選項。 目前,只有當你的功能應用程式託管在 專用(App Service)方案時,才能使用這些屬性。

對於孤立工作者模型,請使用應用程式設定表達式, HttpsCaPem 而非將 PEM 值放入屬性中:

[KafkaTrigger(
    "BrokerList",
    "topic",
    HttpsCaPem = "%KafkaHttpsCaPem%"
)]

在 Azure 中,將KafkaHttpsCaPem應用程式設定為包含 PEM 值的秘密的 金鑰保存庫 參考。 這個 KafkaHttpsCaPem 設定可能像這個例子,其中 <keyVaultName> 是你保險庫的名稱:

@Microsoft.KeyVault(SecretUri=https://<keyVaultName>.vault.azure.net/secrets/httpscapem)

註釋

註 KafkaTrigger 解讓你能建立一個在收到主題時執行的函式。 支援的選項包括下列元素:

元素 描述
name (必要)代表函式程式代碼中佇列或主題訊息的變數名稱。
brokerList (必要)觸發程式所監視的 Kafka 訊息代理程式清單。 如需詳細資訊,請參閱 連線 。
topic (必要)觸發程式所監視的主題。
基數 (選擇性)表示觸發程式輸入的基數。 支援的值為 ONE (預設值) 和 MANY。 ONE當輸入是單一訊息,當MANY輸入是訊息數位時使用。 當您使用 MANY時,也必須設定 dataType。
dataType 定義 Functions 如何處理參數值。 根據預設,會以字串形式取得值,Functions 會嘗試將字串還原串行化為實際的純舊 Java 物件 (POJO)。 當 為 時 string,輸入會視為字串。 當 為 時 binary,訊息會以二進位數據的形式接收,而 Functions 會嘗試將它還原串行化為實際的參數類型 byte[]。
consumerGroup (選擇性)觸發程式所使用的 Kafka 取用者群組。
avroSchema (選擇性)使用 Avro 通訊協定時,一般記錄的架構。
authenticationMode (選擇性)使用簡單驗證和安全性層 (SASL) 驗證時的驗證模式。 支援的值為NotSet(預設值)、Gssapi、PlainScramSha256ScramSha512。
username (選擇性)SASL 驗證的用戶名稱。 當 為AuthenticationMode時Gssapi不受支援。 如需詳細資訊,請參閱 連線 。
password (選擇性)SASL 驗證的密碼。 當 為AuthenticationMode時Gssapi不受支援。 如需詳細資訊,請參閱 連線 。
protocol (選擇性)與訊息代理程式通訊時所使用的安全性通訊協定。 支援的值為NotSet(預設值)、plaintext、sslsasl_plaintextsasl_ssl。
sslCaLocation (選擇性)用於驗證訊息代理程序憑證的 CA 憑證檔案路徑。
sslCertificateLocation (選擇性)用戶端憑證的路徑。
sslKeyLocation (選擇性)用於驗證之用戶端私鑰 (PEM) 的路徑。
sslKeyPassword (選擇性)用戶端憑證的密碼。
延遲閾值 (可選)觸發器的延遲閾值。
schemaRegistryUrl (可選)Avro 架構登錄檔的網址。 如需詳細資訊,請參閱 連線 。
schemaRegistryUsername (可選)Avro 架構登錄檔的用戶名。 如需詳細資訊,請參閱 連線 。
schemaRegistryPassword (可選)Avro 架構登錄檔的密碼。 如需詳細資訊,請參閱 連線 。

組態

下表說明您在 function.json 檔案中設定的繫結設定屬性。

function.json 屬性 描述
type (必修)設定為 kafkaTrigger。
direction (必修)設定為 in。
name (必要)代表函式程式碼中代理數據的變數名稱。
brokerList (必要)觸發程式所監視的 Kafka 訊息代理程式清單。 如需詳細資訊,請參閱 連線 。
topic (必要)觸發程式所監視的主題。
基數 (選擇性)表示觸發程式輸入的基數。 支援的值為 ONE (預設值) 和 MANY。 ONE當輸入是單一訊息,當MANY輸入是訊息數位時使用。 當您使用 MANY時,也必須設定 dataType。
dataType 定義 Functions 如何處理參數值。 根據預設,會以字串形式取得值,Functions 會嘗試將字串還原串行化為實際的純舊 Java 物件 (POJO)。 當 為 時 string,輸入會視為字串。 當 binary時,訊息會以二進位資料形式接收,函式會嘗試將其反序列化為實際的位元組陣列參數型態。
consumerGroup (選擇性)觸發程式所使用的 Kafka 取用者群組。
avroSchema (選擇性)使用 Avro 通訊協定時,一般記錄的架構。
keyAvroSchema (可選)使用 Avro 協定時,訊息金鑰的通用記錄架構。
keyDataType (可選)資料型別以接收訊息金鑰,如同卡夫卡主題。 若 keyAvroSchema 設定為 ,則此值為通用記錄。 公認的值為 Int、 Long、 String和 Binary。
authenticationMode (選擇性)使用簡單驗證和安全性層 (SASL) 驗證時的驗證模式。 支援的值為NotSet(預設值)、Gssapi、PlainScramSha256ScramSha512。
username (選擇性)SASL 驗證的用戶名稱。 當 為AuthenticationMode時Gssapi不受支援。 如需詳細資訊,請參閱 連線 。
password (選擇性)SASL 驗證的密碼。 當 為AuthenticationMode時Gssapi不受支援。 如需詳細資訊,請參閱 連線 。
protocol (選擇性)與訊息代理程式通訊時所使用的安全性通訊協定。 支援的值為NotSet(預設值)、plaintext、sslsasl_plaintextsasl_ssl。
sslCaLocation (選擇性)用於驗證訊息代理程序憑證的 CA 憑證檔案路徑。
sslCertificateLocation (選擇性)用戶端憑證的路徑。
sslKeyLocation (選擇性)用於驗證之用戶端私鑰 (PEM) 的路徑。
sslKeyPassword (選擇性)用戶端憑證的密碼。
sslCertificatePEM (可選)以 PEM 格式呈現的用戶端憑證,作為字串。 如需詳細資訊,請參閱 連線 。
sslKeyPEM (可選)以 PEM 格式呈現的客戶端私鑰,作為字串。 如需詳細資訊,請參閱 連線 。
sslCaPEM (可選)CA 憑證以 PEM 格式以字串形式呈現。 如需詳細資訊,請參閱 連線 。
sslCertificate與KeyPEM (可選)客戶端憑證與金鑰以 PEM 格式呈現字串。 如需詳細資訊,請參閱 連線 。
延遲閾值 (可選)觸發器的延遲閾值。
schemaRegistryUrl (可選)Avro 架構登錄檔的網址。 如需詳細資訊,請參閱 連線 。
schemaRegistryUsername (可選)Avro 架構登錄檔的用戶名。 如需詳細資訊,請參閱 連線 。
schemaRegistryPassword (可選)Avro 架構登錄檔的密碼。 如需詳細資訊,請參閱 連線 。
oAuthBearerMethod (可選)OAuth Bearer 方法。 接受的值是 oidc 和 default。
oAuthBearerClientId (可選)當 oAuthBearerMethod 設為 oidc時,這表示 OAuth 承載的客戶端 ID。 如需詳細資訊,請參閱 連線 。
oAuthBearerClientSecret (可選)當 oAuthBearerMethod 設定為 oidc時,這會指定 OAuth 承載的客戶端秘密。 如需詳細資訊,請參閱 連線 。
oAuthBearerScope(授權承載者範圍) (可選)指定向經紀人申請存取的範圍。
oAuthBearerTokenEndpointUrl (可選)OAuth/OIDC 發行者令牌端點 HTTP(S) URI 用於在使用方法時 oidc 取得令牌。 如需詳細資訊,請參閱 連線 。

組態

下表說明您在 function.json 檔案中設定的繫結設定屬性。 Python 使用snake_case命名規則來描述設定屬性。

function.json 屬性 描述
type (必修)設定為 kafkaTrigger。
direction (必修)設定為 in。
name (必要)代表函式程式碼中代理數據的變數名稱。
broker_list (必要)觸發程式所監視的 Kafka 訊息代理程式清單。 如需詳細資訊,請參閱 連線 。
topic (必要)觸發程式所監視的主題。
基數 (選擇性)表示觸發程式輸入的基數。 支援的值為 ONE (預設值) 和 MANY。 ONE當輸入是單一訊息,當MANY輸入是訊息數位時使用。 當您使用 MANY時,也必須設定 data_type。
data_type 定義 Functions 如何處理參數值。 根據預設,會以字串形式取得值,Functions 會嘗試將字串還原串行化為實際的純舊 Java 物件 (POJO)。 當 為 時 string,輸入會視為字串。 當 為 時 binary,訊息會以二進位數據的形式接收,而 Functions 會嘗試將它還原串行化為實際的參數類型 byte[]。
consumerGroup (選擇性)觸發程式所使用的 Kafka 取用者群組。
avroSchema (選擇性)使用 Avro 通訊協定時,一般記錄的架構。
authentication_mode (選擇性)使用簡單驗證和安全性層 (SASL) 驗證時的驗證模式。 支援的值為NOTSET(預設值)、Gssapi、PlainScramSha256ScramSha512。
username (選擇性)SASL 驗證的用戶名稱。 當 為authentication_mode時Gssapi不受支援。 如需詳細資訊,請參閱 連線 。
password (選擇性)SASL 驗證的密碼。 當 為authentication_mode時Gssapi不受支援。 如需詳細資訊,請參閱 連線 。
protocol (選擇性)與訊息代理程式通訊時所使用的安全性通訊協定。 支援的值為NOTSET(預設值)、plaintext、sslsasl_plaintextsasl_ssl。
sslCaLocation (選擇性)用於驗證訊息代理程序憑證的 CA 憑證檔案路徑。
sslCertificateLocation (選擇性)用戶端憑證的路徑。
sslKeyLocation (選擇性)用於驗證之用戶端私鑰 (PEM) 的路徑。
sslKeyPassword (選擇性)用戶端憑證的密碼。
lag_threshold (可選)觸發器的延遲閾值。
schema_registry_url (可選)Avro 架構登錄檔的網址。 如需詳細資訊,請參閱 連線 。
schema_registry_username (可選)Avro 架構登錄檔的用戶名。 如需詳細資訊,請參閱 連線 。
schema_registry_password (可選)Avro 架構登錄檔的密碼。 如需詳細資訊,請參閱 連線 。
o_auth_bearer_method (可選)OAuth Bearer 方法。 接受的值是 oidc 和 default。
o_auth_bearer_client_id (可選)當 o_auth_bearer_method 設為 oidc時,這表示 OAuth 承載的客戶端 ID。 如需詳細資訊,請參閱 連線 。
o_auth_bearer_client_secret (可選)當 o_auth_bearer_method 設定為 oidc時,這會指定 OAuth 承載的客戶端秘密。 如需詳細資訊,請參閱 連線 。
o_auth_bearer_scope (可選)指定向經紀人申請存取的範圍。
o_auth_bearer_token_endpoint_url (可選)OAuth/OIDC 發行者令牌端點 HTTP(S) URI 用於在使用方法時 oidc 取得令牌。 如需詳細資訊,請參閱 連線 。

注意

憑證相關的PEM屬性和Avro鍵相關的屬性目前還不會出現在Python函式庫中。

使用方式

Kafka 觸發器目前支援以字串和字串陣列形式呈現 Kafka 事件,這些是 JSON 有效載荷。

Kafka 觸發器會將 Kafka 訊息以字串形式傳遞給函式。 此觸發器也支援 JSON 有效載荷的字串陣列。

在高級方案中,你必須啟用 Kafka 輸出的執行時規模監控,才能擴展到多個實例。 若要深入瞭解,請參閱 啟用運行時間調整。

你無法使用 Azure 入口網站 Code + Test 頁面的測試/執行功能來處理 Kafka 觸發器。 您必須改為將測試事件直接傳送至觸發程式所監視的主題。

如需 Kafka 觸發程式支援的一組完整host.json設定,請參閱 host.json設定。

連線

Kafka 綁定擴充功能不支援受管理身份連線。 你必須使用以下其中一種方法來驗證你的 Kafka 連線:

  • 金鑰保存庫 參考:將你的 Kafka 憑證(密碼、API 金鑰、憑證)存入 Azure Key Vault,並從應用程式設定中引用。 你的函式應用程式是透過受管理身份連接到 金鑰保存庫。 欲了解更多資訊,請參閱定義 金鑰保存庫 連接。
  • App Configuration 參考:將連線設定儲存在 Azure 應用程式組態,該設定也可以參考 金鑰保存庫 來存取秘密。 如需詳細資訊,請參閱 Azure 應用程式組態。
  • 共享秘密:將憑證直接儲存在應用程式設定中(靜態時加密)。 欲了解更多資訊,請參閱 定義共享秘密連線。

欲了解更多連線安全資訊,請參閱「管理 Azure Functions 中的連線」。

重要

認證設定必須參考 應用程式設定。 不要在您的程式代碼或組態檔中硬式編碼認證。 在本機執行時,請針對您的認證使用 local.settings.json 檔案 ,而且不會發佈local.settings.json檔案。

當連接到 Confluent 在 Azure 提供的受管理 Kafka 叢集時,你可以使用以下其中一種認證方法。

注意

使用 Flex Consumption 方案時,不支援基於位置的檔案憑證認證屬性(SslCaLocation, SslCertificateLocation, SslKeyLocation)。 相反地,請使用 PEM 基礎的憑證屬性(SslCaPEM、、SslCertificatePEMSslKeyPEM、SslCertificateandKeyPEM)或將憑證存放在 Azure Key Vault。

架構註冊表

要使用 Confluent 在 Kafka 擴充中提供的結構登錄,請設定以下憑證:

設定 建議值 描述
SchemaRegistryUrl SchemaRegistryUrl 用於結構管理的結構登錄服務網址。 通常是格式的 https://psrc-xyz.us-east-2.aws.confluent.cloud
SchemaRegistryUsername CONFLUENT_API_KEY 使用者名稱用於結構登錄檔的基本認證(如有需要)。
SchemaRegistryPassword CONFLUENT_API_SECRET 密碼用於結構登錄檔的基本認證(如有需要)。

使用者名稱/密碼驗證

使用這種認證方式時,請確保 Protocol 設定SaslPlaintextSaslSsl為 或 ,AuthenticationMode被設為 Plain或 ,ScramSha256ScramSha512如果所使用的 CA 憑證與預設的 ISRG Root X1 憑證不同,請確保更新SslCaLocation或 SslCaPEM。

設定 建議值 描述
BrokerList BootstrapServer 名為 BootstrapServer 的應用程式設定包含 Confluent Cloud 設定頁面中找到的啟動程式伺服器值。 值類似於 xyz-xyzxzy.westeurope.azure.confluent.cloud:9092。
使用者名稱 ConfluentCloudUsername 名為 ConfluentCloudUsername 的應用程式設定包含來自 Confluent Cloud 網站的 API 存取金鑰。
密碼 ConfluentCloudPassword 名為 ConfluentCloudPassword 的應用程式設定包含從 Confluent Cloud 網站取得的 API 秘密。
SslCaPEM %SSLCaPemCertificate% 應用程式設定SSLCaPemCertificate,名稱為 PEM 格式的 Azure Key Vault 秘密,包含 CA 憑證。

SSL 認證

請確定 設定 Protocol 為 SSL。

設定 建議值 描述
BrokerList BootstrapServer 名為 BootstrapServer 的應用程式設定包含 Confluent Cloud 設定頁面中找到的啟動程式伺服器值。 值類似於 xyz-xyzxzy.westeurope.azure.confluent.cloud:9092。
SslCaPEM %SslCaCertificatePem% 應用程式設定SslCaCertificatePem,名稱為 PEM 格式的 Azure Key Vault 秘密,包含 CA 憑證。
SslCertificatePEM %SslClientCertificatePem% 應用程式設定,該設定SslClientCertificatePem引用包含 PEM 格式客戶憑證的 Azure Key Vault 秘密。
SslKeyPEM %SslClientKeyPem% 一個名為 SslClientKeyPem App 的設定,引用包含 PEM 格式客戶端私鑰的 Azure Key Vault 秘密。
SslCertificate與KeyPEM %SslClientCertificateAndKeyPem% 應用程式設定,該設定SslClientCertificateAndKeyPem引用包含 Azure Key Vault 秘密,包含以 PEM 格式連接的客戶端憑證與客戶端私鑰。
SslKeyPassword %SslClientKeyPassword% 一個名為 SslClientKeyPassword App 的設定,會引用包含私鑰密碼(如果有的話)的 Azure Key Vault 秘密。

將憑證和私鑰值存放在 Azure Key Vault,而不是直接放在你的函式應用程式設定裡。 將相應的應用程式設定設為 金鑰保存庫 參考,例如:

@Microsoft.KeyVault(SecretUri=https://<keyVaultName>.vault.azure.net/secrets/<secretName>)

OAuth 驗證

使用 OAuth 認證時,請在綁定定義中設定與 OAuth 相關的屬性。

您用於這些設定的字串值必須以應用程式設定的形式出現在 Azure 或Values本機開發期間local.settings.json檔案的集合中。

你也應該在綁定定義中設定 Protocol 和 AuthenticationMode 。

下一步