הצגת אשכול של שירות מנוהל של Google Cloud ל-Apache Kafka

כדי להציג אשכול, אפשר להשתמש במסוף Google Cloud , ב-Google Cloud CLI, בספריית הלקוח או ב-Managed Kafka API. אי אפשר להשתמש ב-API של Apache Kafka בקוד פתוח כדי להציג אשכול.

התפקידים וההרשאות שנדרשים כדי להציג אשכול

כדי לקבל את ההרשאות שנדרשות לצפייה באשכול, צריך לבקש מהאדמין להקצות לכם את תפקיד ה-IAM‏ Managed Kafka Viewer (roles/managedkafka.viewer) בפרויקט. כדי לקרוא הסבר על מתן תפקידים, ראו איך מנהלים את הגישה ברמת הפרויקט, התיקייה והארגון.

זהו תפקיד שמוגדר מראש וכולל את ההרשאות שנדרשות לצפייה באשכול. כדי לראות בדיוק אילו הרשאות נדרשות, אפשר להרחיב את הקטע ההרשאות הנדרשות:

ההרשאות הנדרשות

כדי להציג אשכול, צריך את ההרשאות הבאות:

  • רשימת אשכולות: managedkafka.clusters.list
  • קבלת פרטים על אשכול: managedkafka.clusters.get

יכול להיות שתקבלו את ההרשאות האלה באמצעות תפקידים בהתאמה אישית או תפקידים מוגדרים מראש אחרים.

הצגת אשכול

כדי לראות אשכול, מבצעים את השלבים הבאים:

המסוף

  1. נכנסים לדף Clusters במסוף Google Cloud .

    מעבר אל Clusters

    מוצגת רשימה של האשכולות שיצרתם בפרויקט.

  2. כדי לראות אשכול ספציפי, לוחצים על שם האשכול.

  3. ייפתח דף הפרטים של האשכול. בדף הזה יש את הכרטיסיות הבאות:

    • משאבים: מוצגת רשימה של נושאים וקבוצות צרכנים שמשויכים לאשכול.
    • בקרת גישה: מוצגים רשומות ה-ACL של האשכול.
    • הגדרות: מציג את ההגדרה של האשכול, כולל רשימת רשתות המשנה שמשויכות לאשכול.
    • מעקב: כאן מוצגות ההתראות על מעקב שמשויכות לאשכול.
    • יומנים: מוצגים היומנים שקשורים לאשכולות מ-Logs Explorer.
    • מקורות: קיצורי דרך ליצירת מחברי מקורות של Kafka Connect או ליצירת נתוני בדיקה סינתטיים.
    • Sinks: מספק קיצורי דרך ליצירת מחברי Kafka Connect sink.

gcloud

  1. במסוף Google Cloud , מפעילים את Cloud Shell.

    הפעלת Cloud Shell

    בחלק התחתון של Google Cloud המסוף יתחיל סשן של Cloud Shell ותופיע הודעה של שורת הפקודה. Cloud Shell היא סביבת מעטפת שבה ה-CLI של Google Cloud מותקן ומוגדרים ערכים לפרויקט הקיים. הסשן יופעל תוך כמה שניות.

  2. לפני השימוש בנתוני הפקודה הבאים, צריך להחליף את הנתונים הבאים:

    • ‫PROJECT_ID: מזהה הפרויקט.
    • ‫LOCATION: המיקום של האשכול.
    • ‫CLUSTER_ID: מזהה האשכול.

    מריצים את הפקודה הבאה:

    ‫Linux,‏ macOS או Cloud Shell

    gcloud managed-kafka clusters describe CLUSTER_ID \
        --location=LOCATION \
        --project=PROJECT_ID

    ‏Windows (PowerShell)

    gcloud managed-kafka clusters describe CLUSTER_ID `
        --location=LOCATION `
        --project=PROJECT_ID

    Windows‏ (cmd.exe)

    gcloud managed-kafka clusters describe CLUSTER_ID ^
        --location=LOCATION ^
        --project=PROJECT_ID

    אמורים לקבל תגובה שדומה לזו:

    bootstrapAddress: bootstrap.CLUSTER_ID.LOCATION.managedkafka.PROJECT_ID.cloud.goog:9092
    bootstrapAddressMTLS: bootstrap.CLUSTER_ID.LOCATION.managedkafka.PROJECT_ID.cloud.goog:9192
    brokerCapacityConfig: {}
    capacityConfig:
      memoryBytes: 'MEMORY'
      vcpuCount: 'CPU_COUNT'
    createTime: 'CREATE_TIME'
    gcpConfig:
      accessConfig:
        networkConfigs:
        - subnet: projects/PROJECT_ID/regions/LOCATION/subnetworks/SUBNET_ID
        publicClusterConfig:
          allowedSourceIpRanges:
          - 203.0.113.5/32
    labels:
      goog-terraform-provisioned: 'true'
    name: projects/PROJECT_ID/locations/LOCATION/clusters/CLUSTER_ID
    publicClusterDetails:
      discoveryDnsRecords:
      - discovery-1.CLUSTER_ID.LOCATION.managedkafka.PROJECT_ID.cloud.goog
      externalIpAddresses:
      - 203.0.113.1
      - 203.0.113.2
      - 203.0.113.3
      - 203.0.113.4
    rebalanceConfig:
      mode: AUTO_REBALANCE_ON_SCALE_UP
    satisfiesPzi: false
    satisfiesPzs: false
    state: ACTIVE
    tlsConfig:
      trustConfig: {}
    updateOptions: {}
    updateTime: 'UPDATE_TIME'
    

REST

לפני שמשתמשים בנתוני הבקשה, צריך להחליף את הנתונים הבאים:

  • ‫PROJECT_ID: מזהה הפרויקט ב- Google Cloud
  • ‫LOCATION: המיקום של האשכול.
  • ‫CLUSTER_ID: מזהה האשכול.
  • ‫CLUSTER_VIEW: כמות המטא-נתונים שרוצים להחזיר. מציינים אחד מהערכים הבאים:

    • ‫CLUSTER_VIEW_BASIC: החזרת המטא-נתונים הבסיסיים של האשכול.
    • ‫CLUSTER_VIEW_FULL: מחזירה את כל המטא-נתונים של האשכול, כולל מידע על הברוקרים של האשכול ומספר הגרסה של Kafka שמופעלת באשכול.

    אם לא מציינים שיטה, ברירת המחדל היא CLUSTER_VIEW_BASIC.

ה-method של ה-HTTP וכתובת ה-URL:

GET https://managedkafka.googleapis.com/v1/projects/PROJECT_ID/locations/LOCATION/clusters/CLUSTER_ID?view=CLUSTER_VIEW

כדי לשלוח את הבקשה צריך להרחיב אחת מהאפשרויות הבאות:

אתם אמורים לקבל תגובת JSON שדומה לזו:

{
  "name": "projects/PROJECT_ID/locations/LOCATION/clusters/CLUSTER_ID",
  "createTime": "CREATE_TIME",
  "updateTime": "UPDATE_TIME",
  "capacityConfig": {
    "vcpuCount": "CPU_COUNT",
    "memoryBytes": "MEMORY"
  },
  "rebalanceConfig": {},
  "gcpConfig": {
    "accessConfig": {
      "networkConfigs": [
        {
          "subnet": "projects/PROJECT_ID/locations/LOCATION/subnetworks/SUBNET_ID"
        }
      ],
      "publicClusterConfig": {
        "allowedSourceIpRanges": [
          "203.0.113.5/32"
        ]
      }
    }
  },
  "state": "ACTIVE",
  "satisfiesPzi": false,
  "satisfiesPzs": false,
  "tlsConfig": {
    "trustConfig": {}
  },
  "updateOptions": {},
  "publicClusterDetails": {
    "discoveryDnsRecords": [
      "discovery-1.CLUSTER_ID.LOCATION.managedkafka.PROJECT_ID.cloud.goog"
    ],
    "externalIpAddresses": [
      "203.0.113.1",
      "203.0.113.2",
      "203.0.113.3",
      "203.0.113.4"
    ]
  }
}

המשך

לפני שמנסים את הדוגמה הזו, צריך לפעול לפי הוראות ההגדרה של Go בקטע התקנת ספריות הלקוח. מידע נוסף מופיע ב מאמרי העזרה של ה-API של Go עבור שירות מנוהל ל-Apache Kafka.

כדי לבצע אימות לשירות המנוהל ל-Apache Kafka, צריך להגדיר את Application Default Credentials‏(ADC). מידע נוסף זמין במאמר הגדרת ADC לסביבת פיתוח מקומית.

import (
	"context"
	"fmt"
	"io"

	"cloud.google.com/go/managedkafka/apiv1/managedkafkapb"
	"google.golang.org/api/option"

	managedkafka "cloud.google.com/go/managedkafka/apiv1"
)

func getCluster(w io.Writer, projectID, region, clusterID string, opts ...option.ClientOption) error {
	// projectID := "my-project-id"
	// region := "us-central1"
	// clusterID := "my-cluster"
	ctx := context.Background()
	client, err := managedkafka.NewClient(ctx, opts...)
	if err != nil {
		return fmt.Errorf("managedkafka.NewClient got err: %w", err)
	}
	defer client.Close()

	clusterPath := fmt.Sprintf("projects/%s/locations/%s/clusters/%s", projectID, region, clusterID)
	req := &managedkafkapb.GetClusterRequest{
		Name: clusterPath,
	}
	cluster, err := client.GetCluster(ctx, req)
	if err != nil {
		return fmt.Errorf("client.GetCluster got err: %w", err)
	}
	fmt.Fprintf(w, "Got cluster: %#v\n", cluster)
	return nil
}

Java

לפני שמנסים את הדוגמה הזו, צריך לפעול לפי הוראות ההגדרה של Java במאמר התקנת ספריות הלקוח. מידע נוסף מופיע ב מאמרי העזרה של Managed Service for Apache Kafka Java API.

כדי לבצע אימות לשירות המנוהל ל-Apache Kafka, מגדירים את ה-Application Default Credentials. מידע נוסף זמין במאמר הגדרת ADC לסביבת פיתוח מקומית.

import com.google.api.gax.rpc.ApiException;
import com.google.cloud.managedkafka.v1.Cluster;
import com.google.cloud.managedkafka.v1.ClusterName;
import com.google.cloud.managedkafka.v1.ManagedKafkaClient;
import java.io.IOException;

public class GetCluster {

  public static void main(String[] args) throws Exception {
    // TODO(developer): Replace these variables before running the example.
    String projectId = "my-project-id";
    String region = "my-region"; // e.g. us-east1
    String clusterId = "my-cluster";
    getCluster(projectId, region, clusterId);
  }

  public static void getCluster(String projectId, String region, String clusterId)
      throws Exception {
    try (ManagedKafkaClient managedKafkaClient = ManagedKafkaClient.create()) {
      // This operation is being handled synchronously.
      Cluster cluster = managedKafkaClient.getCluster(ClusterName.of(projectId, region, clusterId));
      System.out.println(cluster.getAllFields());
    } catch (IOException | ApiException e) {
      System.err.printf("managedKafkaClient.getCluster got err: %s", e.getMessage());
    }
  }
}

Python

לפני שמנסים את הדוגמה הזו, פועלים לפי הוראות ההגדרה של Python במאמר התקנת ספריות הלקוח. מידע נוסף מופיע ב מאמרי העזרה של ה-API בשפת Python של שירות מנוהל ל-Apache Kafka.

כדי לבצע אימות לשירות המנוהל ל-Apache Kafka, מגדירים את ה-Application Default Credentials. מידע נוסף זמין במאמר הגדרת ADC לסביבת פיתוח מקומית.

from google.api_core.exceptions import NotFound
from google.cloud import managedkafka_v1

# TODO(developer)
# project_id = "my-project-id"
# region = "us-central1"
# cluster_id = "my-cluster"

client = managedkafka_v1.ManagedKafkaClient()

cluster_path = client.cluster_path(project_id, region, cluster_id)
request = managedkafka_v1.GetClusterRequest(
    name=cluster_path,
)

try:
    cluster = client.get_cluster(request=request)
    print("Got cluster:", cluster)
except NotFound as e:
    print(f"Failed to get cluster {cluster_id} with error: {e.message}")

הצגת המאפיינים והמשאבים של אשכול

בקטעים הבאים מוסבר איך לקבל פרטים על מאפיינים ומשאבים שונים שמשויכים לאשכול של שירות מנוהל של Google Cloud ל-Apache Kafka.

כתובת אתחול

לקוחות Kafka משתמשים בכתובת ה-bootstrap של האשכול כדי ליצור חיבור לאשכול. כתובת ה-bootstrap קבועה למשך חיי האשכול, אבל הפורמט של כתובת URL התחלתית עשוי להיות שונה בין אשכולות.

אם הופעלה גישה ציבורית לאשכול, כתובת ה-bootstrap מומרת אוטומטית לכתובות IP ציבוריות או פרטיות, בהתאם למקום שממנו הלקוח מתחבר (באמצעות DNS עם אופק מפוצל). גם לקוחות פנימיים וגם לקוחות חיצוניים צריכים להשתמש באותה כתובת bootstrap כדי להתחבר לאשכול. אל תגדירו את לקוחות Kafka להתחבר לרשומות DNS של discovery, ששמורות רק להגדרת חומות אש חיצוניות של תעבורת נתונים יוצאת.

כדי לקבל את כתובת ה-bootstrap:

המסוף

  1. עוברים לדף שירות מנוהל ל-Apache Kafka > Clusters (אשכולות).

    מעבר אל Clusters

  2. לוחצים על שם האשכול.

  3. בוחרים בכרטיסייה Configurations.

  4. אם אתם משתמשים ב-SASL לאימות, כתובת URL התחלתית מופיעה בקטע Bootstrap URL.

    אם משתמשים ב-mutual TLS (mTLS) לאימות, כתובת URL התחלתית מופיעה בקטע mTLS Bootstrap URL.

    לוחצים על העתקה כדי להעתיק את הערך.

gcloud

כדי לקבל את כתובת URL התחלתית, משתמשים בפקודה managed-kafka clusters describe.

אם אתם משתמשים ב-SASL לאימות, מריצים את הפקודה הבאה:

gcloud managed-kafka clusters describe CLUSTER_ID \
    --location=LOCATION \
    --format="value(bootstrapAddress)"

אם אתם משתמשים ב-mutual TLS (mTLS) כדי לבצע אימות, מריצים את הפקודה הבאה:

gcloud managed-kafka clusters describe CLUSTER_ID \
    --location=LOCATION \
    --format="value(bootstrapAddressMTLS)"

מחליפים את מה שכתוב בשדות הבאים:

  • ‫CLUSTER_ID: המזהה או השם של האשכול.
  • ‫LOCATION: המיקום של האשכול.

מידע נוסף על אימות SASL ו-mTLS זמין במאמר סוגי אימות לשרתי Kafka.

סוכנים

כדי לראות את השרתים המתווכים באשכול, אפשר לעיין במאמר הצגת שרתים מתווכים באשכול של שירות מנוהל ל-Apache Kafka.

קבוצות צרכנים

קבוצת צרכנים היא קבוצה של צרכנים שמשתפים פעולה כדי לצרוך נתונים מנושאים שונים. כדי לראות את קבוצות הצרכנים של אשכול, אפשר לעיין בדפים הבאים:

רשתות משנה

כדי לראות את רשתות המשנה שמחוברות לאשכול, מבצעים את השלבים הבאים:

המסוף

  1. נכנסים לדף Clusters במסוף Google Cloud .

    מעבר אל Clusters

  2. לוחצים על שם האשכול.

  3. בדף הפרטים של האשכול, לוחצים על הכרטיסייה Configurations. רשתות משנה מקושרות מפורטות בקטע Subnetworks.

gcloud

  1. מריצים את הפקודה gcloud managed-kafka clusters describe:

    gcloud managed-kafka clusters describe CLUSTER_ID \
        --location=LOCATION \
        --format="yaml(gcpConfig.accessConfig.networkConfigs)"
    

    מחליפים את מה שכתוב בשדות הבאים:

    • ‫CLUSTER_ID: המזהה או השם של האשכול.
    • ‫LOCATION: המיקום של האשכול.

מידע נוסף זמין במאמר בנושא הגדרת רשתות בשירות המנוהל ל-Apache Kafka.

טווחים מותרים של כתובות IP של מקורות

כדי לראות את טווחי כתובות ה-IP של המקור שמותרים לאשכול ציבורי, מבצעים את השלבים הבאים:

המסוף

  1. נכנסים לדף Clusters במסוף Google Cloud .

    מעבר אל Clusters

  2. לוחצים על שם האשכול.

  3. לוחצים על הכרטיסייה Configurations. טווחי כתובות ה-IP המותרים של המקור מפורטים בטבלה Allowed source IP ranges (טווחי כתובות ה-IP המותרים של המקור) בעמודה IP range (טווח כתובות ה-IP).

gcloud

  1. כדי להציג את טווחי כתובות ה-IP המותרים של המקור באמצעות gcloud, בודקים את השדה gcpConfig.accessConfig.publicClusterConfig.allowedSourceIpRanges בתשובה של הפקודה gcloud managed-kafka clusters describe:

    gcloud managed-kafka clusters describe CLUSTER_ID \
        --location=LOCATION \
        --format="yaml(gcpConfig.accessConfig.publicClusterConfig.allowedSourceIpRanges)"
    

    מחליפים את מה שכתוב בשדות הבאים:

    • ‫CLUSTER_ID: המזהה או השם של האשכול.
    • ‫LOCATION: המיקום של האשכול.

REST

כדי לראות את טווחי כתובות ה-IP המותרים של המקור באמצעות API בארכיטקטורת REST, בודקים את השדה gcpConfig.accessConfig.publicClusterConfig.allowedSourceIpRanges בתשובה של בקשת GET:

לפני שמשתמשים בנתוני הבקשה, צריך להחליף את הנתונים הבאים:

  • ‫PROJECT_ID: מזהה הפרויקט ב- Google Cloud
  • ‫LOCATION: המיקום של האשכול.
  • ‫CLUSTER_ID: מזהה האשכול.
  • ‫CLUSTER_VIEW: כמות המטא-נתונים שרוצים להחזיר. מציינים אחד מהערכים הבאים:

    • ‫CLUSTER_VIEW_BASIC: החזרת המטא-נתונים הבסיסיים של האשכול.
    • ‫CLUSTER_VIEW_FULL: מחזירה את כל המטא-נתונים של האשכול, כולל מידע על הברוקרים של האשכול ומספר הגרסה של Kafka שמופעלת באשכול.

    אם לא מציינים שיטה, ברירת המחדל היא CLUSTER_VIEW_BASIC.

ה-method של ה-HTTP וכתובת ה-URL:

GET https://managedkafka.googleapis.com/v1/projects/PROJECT_ID/locations/LOCATION/clusters/CLUSTER_ID?view=CLUSTER_VIEW

כדי לשלוח את הבקשה צריך להרחיב אחת מהאפשרויות הבאות:

אתם אמורים לקבל תגובת JSON שדומה לזו:

{
  "name": "projects/PROJECT_ID/locations/LOCATION/clusters/CLUSTER_ID",
  "createTime": "CREATE_TIME",
  "updateTime": "UPDATE_TIME",
  "capacityConfig": {
    "vcpuCount": "CPU_COUNT",
    "memoryBytes": "MEMORY"
  },
  "rebalanceConfig": {},
  "gcpConfig": {
    "accessConfig": {
      "networkConfigs": [
        {
          "subnet": "projects/PROJECT_ID/locations/LOCATION/subnetworks/SUBNET_ID"
        }
      ],
      "publicClusterConfig": {
        "allowedSourceIpRanges": [
          "203.0.113.5/32"
        ]
      }
    }
  },
  "state": "ACTIVE",
  "satisfiesPzi": false,
  "satisfiesPzs": false,
  "tlsConfig": {
    "trustConfig": {}
  },
  "updateOptions": {},
  "publicClusterDetails": {
    "discoveryDnsRecords": [
      "discovery-1.CLUSTER_ID.LOCATION.managedkafka.PROJECT_ID.cloud.goog"
    ],
    "externalIpAddresses": [
      "203.0.113.1",
      "203.0.113.2",
      "203.0.113.3",
      "203.0.113.4"
    ]
  }
}

מידע נוסף על טווחי כתובות IP מותרים של מקורות זמין במאמר בנושא אשכולות ציבוריים.

פרטי האשכול הציבורי

כדי לראות את פרטי האשכול הציבורי, כולל כתובות ה-IP החיצוניות ורשומות ה-DNS של הגילוי, מבצעים את השלבים הבאים:

המסוף

  1. נכנסים לדף Clusters במסוף Google Cloud .

    מעבר אל Clusters

  2. לוחצים על שם האשכול.

  3. לוחצים על הכרטיסייה Configurations. כתובות ה-IP החיצוניות ורשומות ה-DNS של הגילוי מופיעות בשדות External IP Addresses ו-DNS Discovery Records.

gcloud

  1. מריצים את הפקודה gcloud managed-kafka clusters describe:

    gcloud managed-kafka clusters describe CLUSTER_ID \
        --location=LOCATION \
        --format="yaml(publicClusterDetails)"
    

    מחליפים את מה שכתוב בשדות הבאים:

    • ‫CLUSTER_ID: המזהה או השם של האשכול.
    • ‫LOCATION: המיקום של האשכול.

REST

כדי לראות את הפרטים של האשכול הציבורי באמצעות API בארכיטקטורת REST, בודקים את השדה publicClusterDetails בתגובה לבקשת GET:

לפני שמשתמשים בנתוני הבקשה, צריך להחליף את הנתונים הבאים:

  • ‫PROJECT_ID: מזהה הפרויקט ב- Google Cloud
  • ‫LOCATION: המיקום של האשכול.
  • ‫CLUSTER_ID: מזהה האשכול.
  • ‫CLUSTER_VIEW: כמות המטא-נתונים שרוצים להחזיר. מציינים אחד מהערכים הבאים:

    • ‫CLUSTER_VIEW_BASIC: החזרת המטא-נתונים הבסיסיים של האשכול.
    • ‫CLUSTER_VIEW_FULL: מחזירה את כל המטא-נתונים של האשכול, כולל מידע על הברוקרים של האשכול ומספר הגרסה של Kafka שמופעלת באשכול.

    אם לא מציינים שיטה, ברירת המחדל היא CLUSTER_VIEW_BASIC.

ה-method של ה-HTTP וכתובת ה-URL:

GET https://managedkafka.googleapis.com/v1/projects/PROJECT_ID/locations/LOCATION/clusters/CLUSTER_ID?view=CLUSTER_VIEW

כדי לשלוח את הבקשה צריך להרחיב אחת מהאפשרויות הבאות:

אתם אמורים לקבל תגובת JSON שדומה לזו:

{
  "name": "projects/PROJECT_ID/locations/LOCATION/clusters/CLUSTER_ID",
  "createTime": "CREATE_TIME",
  "updateTime": "UPDATE_TIME",
  "capacityConfig": {
    "vcpuCount": "CPU_COUNT",
    "memoryBytes": "MEMORY"
  },
  "rebalanceConfig": {},
  "gcpConfig": {
    "accessConfig": {
      "networkConfigs": [
        {
          "subnet": "projects/PROJECT_ID/locations/LOCATION/subnetworks/SUBNET_ID"
        }
      ],
      "publicClusterConfig": {
        "allowedSourceIpRanges": [
          "203.0.113.5/32"
        ]
      }
    }
  },
  "state": "ACTIVE",
  "satisfiesPzi": false,
  "satisfiesPzs": false,
  "tlsConfig": {
    "trustConfig": {}
  },
  "updateOptions": {},
  "publicClusterDetails": {
    "discoveryDnsRecords": [
      "discovery-1.CLUSTER_ID.LOCATION.managedkafka.PROJECT_ID.cloud.goog"
    ],
    "externalIpAddresses": [
      "203.0.113.1",
      "203.0.113.2",
      "203.0.113.3",
      "203.0.113.4"
    ]
  }
}

נושאים

כדי לראות את הנושאים באשכול, אפשר לעיין בדפים הבאים:

מה השלב הבא?

‫Apache Kafka®‎ הוא סימן מסחרי רשום של The Apache Software Foundation או של השותפים העצמאיים שלה בארצות הברית או במדינות אחרות.