Clientseitige Messwerte zur Fehlerbehebung bei hoher Latenz verwenden

Memorystore for Redis bietet zwar serverseitige Echtzeitmesswerte zur Überwachung von Durchsatz, CPU-Auslastung und Arbeitsspeichernutzung, aber diese Daten allein erklären möglicherweise nicht, warum in Ihrer Clientanwendung in komplexen verteilten Systemen eine hohe Latenz auftritt.

Clientseitige Messwerte schaffen hier Abhilfe, indem sie Transparenz in den gesamten Anfrage-Antwort-Zyklus bringen. Sie messen einen Befehl ab dem Zeitpunkt, an dem die Anwendung ihn initiiert, bis die Anwendung die Antwort verarbeitet. Durch das Erfassen dieser Datenpunkte können Sie genau bestimmen, ob die Latenz von der Anwendungslogik, dem Netzwerkpfad oder dem Redis-Server stammt.

Hinweis

Achten Sie darauf, dass Ihre Clientanwendung ein Dienstkonto verwendet und ihr die folgenden IAM-Rollen (Identity and Access Management) zugewiesen sind:

  • roles/cloudtrace.agent (Cloud Trace-Agent)
  • roles/monitoring.metricWriter (Monitoring-Messwert-Autor)

Weitere Informationen zum Zuweisen von Rollen finden Sie in der Kurzanleitung IAM-Rolle über die Google Cloud Console zuweisen.

Cloud Monitoring API aktivieren

Wenn Sie clientseitige Messwerte nach Monitoring exportieren möchten, muss die Monitoring API für Ihre Anwendung aktiviert sein. Durch das Exportieren und Visualisieren dieser Messwerte in Monitoring können Sie die Ursache von Engpässen ermitteln und feststellen, wo die Latenz entsteht.

So aktivieren Sie die Monitoring API:

  1. Rufen Sie in der Google Cloud Console die Seite APIs und Dienste auf.

    Zu APIs und Dienste

  2. Wählen Sie das Projekt aus, in dem Sie die Memorystore for Redis-Instanz erstellt haben.

  3. Klicken Sie auf APIs und Dienste aktivieren.

  4. Suchen Sie nach monitoring.

  5. Klicken Sie in den Suchergebnissen auf Cloud Monitoring API.

  6. Wenn API aktiviert angezeigt wird, ist die API bereits aktiviert. Klicken Sie andernfalls auf Aktivieren.

Cloud Trace API aktivieren

Wenn Sie verteilte Traces in Trace ansehen möchten, müssen Sie die Trace API aktivieren. Anschließend können Sie mit dem Trace Explorer diese Traces ansehen, Engpässe diagnostizieren und die Quelle der Latenz in Ihrer Anwendung isolieren.

So aktivieren Sie die Trace API:

  1. Rufen Sie in der Google Cloud Console die Seite APIs und Dienste auf.

    Zu APIs und Dienste

  2. Wählen Sie das Projekt aus, in dem Sie die Memorystore for Redis-Instanz erstellt haben.

  3. Klicken Sie auf APIs und Dienste aktivieren.

  4. Suchen Sie nach trace.

  5. Klicken Sie in den Suchergebnissen auf Cloud Trace API.

  6. Wenn API aktiviert angezeigt wird, ist die API bereits aktiviert. Klicken Sie andernfalls auf Aktivieren.

Clientseitige Messwerte aktivieren

Wenn Sie clientseitige Messwerte aktivieren möchten, fügen Sie dem Code Ihrer Anwendung das OpenTelemetry SDK, den Cloud Monitoring-Exporter und den Cloud Trace-Exporter hinzu. Die OpenTelemetry-Instrumentierung, die direkt in der Redis-Clientbibliothek Ihrer Anwendung ausgeführt wird, erfasst die Messwerte. So kann Ihre Anwendung Latenzdatenpunkte aufzeichnen und zur Visualisierung nach Monitoring und Trace exportieren.

Sie können clientseitige Messwerte in Go, Java, Node.js, oder Python aktivieren. Informationen zum Aktivieren der Messwerte für jede Sprache finden Sie auf den folgenden Tabs.

Go

  1. Führen Sie die folgenden Befehle in Ihrem Terminal aus, um die erforderlichen OpenTelemetry- und Google Cloud Exporter Abhängigkeiten zu installieren:

      go get github.com/gomodule/redigo/redis@latest
      go get go.opentelemetry.io/otel
      go get go.opentelemetry.io/otel/sdk/trace
      go get go.opentelemetry.io/otel/sdk/metric
      go get github.com/GoogleCloudPlatform/opentelemetry-operations-go/exporter/trace
      go get github.com/GoogleCloudPlatform/opentelemetry-operations-go/exporter/metric
  2. Erstellen Sie eine main.go-Datei und fügen Sie ihr den folgenden Code hinzu, um die clientseitigen Messwerte zu aktivieren:

    package main
    
    import (
    	"context"
    	"fmt"
    	"log"
    	"os"
    	"time"
    
    	"github.com/gomodule/redigo/redis"
    	"go.opentelemetry.io/otel"
    	"go.opentelemetry.io/otel/attribute"
    	"go.opentelemetry.io/otel/codes"
    	"go.opentelemetry.io/otel/metric"
    	"go.opentelemetry.io/otel/trace"
    
    	gcpmetric "github.com/GoogleCloudPlatform/opentelemetry-operations-go/exporter/metric"
    	gcptrace "github.com/GoogleCloudPlatform/opentelemetry-operations-go/exporter/trace"
    	sdkmetric "go.opentelemetry.io/otel/sdk/metric"
    	sdktrace "go.opentelemetry.io/otel/sdk/trace"
    )
    
    // MetricClient encapsulates the tracer and metric histograms to avoid package-level globals.
    type MetricClient struct {
    	tracer           trace.Tracer
    	rttHist          metric.Float64Histogram
    	clientBlockHist  metric.Float64Histogram
    	appBlockHist     metric.Float64Histogram
    	retryCounter     metric.Int64Counter
    	connErrorCounter metric.Int64Counter
    }
    
    // sleep hook enables lightning-fast unit tests by stubbing out real time.Sleep
    var sleep = time.Sleep
    
    // sinceMs calculates elapsed time in fractional milliseconds to avoid truncating sub-millisecond durations.
    func sinceMs(start time.Time) float64 {
    	return float64(time.Since(start).Microseconds()) / 1000.0
    }
    
    func initTelemetry(ctx context.Context) (*MetricClient, func(), error) {
    	traceExporter, err := gcptrace.New()
    	if err != nil {
    		return nil, nil, fmt.Errorf("gcptrace.New: %w", err)
    	}
    	tp := sdktrace.NewTracerProvider(sdktrace.WithBatcher(traceExporter))
    	otel.SetTracerProvider(tp)
    	tracer := tp.Tracer("redigo.client")
    
    	metricExporter, err := gcpmetric.New()
    	if err != nil {
    		return nil, nil, fmt.Errorf("gcpmetric.New: %w", err)
    	}
    	mp := sdkmetric.NewMeterProvider(sdkmetric.WithReader(sdkmetric.NewPeriodicReader(metricExporter, sdkmetric.WithInterval(10*time.Second))))
    	otel.SetMeterProvider(mp)
    	meter := mp.Meter("redigo.metrics")
    
    	rttHist, err := meter.Float64Histogram("redis_client_rtt", metric.WithUnit("ms"))
    	if err != nil {
    		return nil, nil, fmt.Errorf("redis_client_rtt histogram: %w", err)
    	}
    	clientBlockHist, err := meter.Float64Histogram("redis_client_blocking_latency", metric.WithUnit("ms"))
    	if err != nil {
    		return nil, nil, fmt.Errorf("redis_client_blocking_latency histogram: %w", err)
    	}
    	appBlockHist, err := meter.Float64Histogram("redis_application_blocking_latency", metric.WithUnit("ms"))
    	if err != nil {
    		return nil, nil, fmt.Errorf("redis_application_blocking_latency histogram: %w", err)
    	}
    	retryCounter, err := meter.Int64Counter("redis_retry_count")
    	if err != nil {
    		return nil, nil, fmt.Errorf("redis_retry_count counter: %w", err)
    	}
    	connErrorCounter, err := meter.Int64Counter("redis_connectivity_error_count")
    	if err != nil {
    		return nil, nil, fmt.Errorf("redis_connectivity_error_count counter: %w", err)
    	}
    
    	client := &MetricClient{
    		tracer:           tracer,
    		rttHist:          rttHist,
    		clientBlockHist:  clientBlockHist,
    		appBlockHist:     appBlockHist,
    		retryCounter:     retryCounter,
    		connErrorCounter: connErrorCounter,
    	}
    
    	initAttrs := metric.WithAttributes(attribute.String("operation", "startup"))
    	client.retryCounter.Add(ctx, 0, initAttrs)
    	client.connErrorCounter.Add(ctx, 0, initAttrs)
    
    	shutdown := func() {
    		tp.Shutdown(ctx)
    		mp.Shutdown(ctx)
    	}
    
    	return client, shutdown, nil
    }
    
    func (c *MetricClient) smartRedisCall(ctx context.Context, pool *redis.Pool, operationName string, commandName string, args ...interface{}) (interface{}, error) {
    	// Create a dedicated child span for the Redis command
    	ctx, span := c.tracer.Start(ctx, operationName)
    	span.SetAttributes(attribute.String("redis.command", commandName))
    	defer span.End()
    
    	maxRetries := 3
    	attempt := 0
    	metricOpts := metric.WithAttributes(attribute.String("operation", operationName))
    	var lastErr error
    
    	for attempt < maxRetries {
    		poolStart := time.Now()
    		// Use GetContext to respect context deadlines and cancellation
    		conn, err := pool.GetContext(ctx)
    		c.clientBlockHist.Record(ctx, sinceMs(poolStart), metricOpts)
    
    		if err != nil {
    			c.connErrorCounter.Add(ctx, 1, metricOpts)
    			c.retryCounter.Add(ctx, 1, metricOpts)
    			span.RecordError(err)
    			span.SetStatus(codes.Error, err.Error())
    			lastErr = err
    			attempt++
    			if attempt >= maxRetries {
    				break
    			}
    			sleep(time.Duration(100<<attempt) * time.Millisecond)
    			continue
    		}
    
    		// Check if the connection is dead
    		if err := conn.Err(); err != nil {
    			conn.Close()
    			c.connErrorCounter.Add(ctx, 1, metricOpts)
    			c.retryCounter.Add(ctx, 1, metricOpts)
    			span.RecordError(err)
    			span.SetStatus(codes.Error, err.Error())
    			lastErr = err
    			attempt++
    			if attempt >= maxRetries {
    				break
    			}
    			sleep(time.Duration(100<<attempt) * time.Millisecond)
    			continue
    		}
    
    		reqStart := time.Now()
    		// Redigo has no native DoContext; pass timeouts using redis.DoWithTimeout when context has a deadline
    		var reply interface{}
    		if deadline, ok := ctx.Deadline(); ok {
    			reply, err = redis.DoWithTimeout(conn, time.Until(deadline), commandName, args...)
    		} else {
    			reply, err = conn.Do(commandName, args...)
    		}
    		c.rttHist.Record(ctx, sinceMs(reqStart), metricOpts)
    		conn.Close()
    
    		if err != nil {
    			c.retryCounter.Add(ctx, 1, metricOpts)
    			span.RecordError(err)
    			span.SetStatus(codes.Error, err.Error())
    			lastErr = err
    			attempt++
    			if attempt >= maxRetries {
    				break
    			}
    			sleep(time.Duration(100<<attempt) * time.Millisecond)
    			continue
    		}
    
    		appStart := time.Now()
    		// Replace fmt.Sprintf to remove unnecessary string formatting overhead
    		sleep(2 * time.Millisecond)
    		c.appBlockHist.Record(ctx, sinceMs(appStart), metricOpts)
    
    		// Reset span status to Ok if the retry or execution eventually succeeds
    		span.SetStatus(codes.Ok, "")
    
    		return reply, nil
    	}
    	return nil, fmt.Errorf("max retries reached for %s: %w", operationName, lastErr)
    }
    
    func main() {
    	ctx := context.Background()
    	client, shutdown, err := initTelemetry(ctx)
    	if err != nil {
    		log.Printf("Failed to initialize telemetry: %v", err)
    		os.Exit(1)
    	}
    	defer shutdown()
    
    	redisHost := os.Getenv("REDISHOST")
    	redisPort := os.Getenv("REDISPORT")
    	if redisPort == "" {
    		redisPort = "6379"
    	}
    
    	pool := &redis.Pool{
    		MaxIdle:     10,
    		MaxActive:   20,
    		IdleTimeout: 240 * time.Second,
    		Wait:        true,
    		Dial: func() (redis.Conn, error) {
    			return redis.Dial("tcp", fmt.Sprintf("%s:%s", redisHost, redisPort))
    		},
    	}
    	defer pool.Close()
    
    	ctx, span := client.tracer.Start(ctx, "fetch_data_span")
    	defer span.End()
    
    	// Simple write and read operations
    	_, err = client.smartRedisCall(ctx, pool, "set_user", "SET", "user:123", "active")
    	if err != nil {
    		log.Printf("Error setting data: %v", err)
    	}
    	val, err := client.smartRedisCall(ctx, pool, "get_user", "GET", "user:123")
    	if err != nil {
    		log.Printf("Error fetching data: %v", err)
    	} else {
    		log.Printf("Retrieved value: %s", val)
    	}
    }
    
  3. Führen Sie Ihre Anwendung mindestens eine Minute lang aus, damit der Exporter genügend Zeit hat, die veröffentlichten Messwerte zu verarbeiten und an Monitoring zu senden.

Java

  1. Fügen Sie der Datei pom.xml Ihrer Anwendung den folgenden Code hinzu, um die erforderlichen OpenTelemetry- und Google Cloud Exporter Abhängigkeiten zu installieren:

    <dependencies>
        <dependency>
            <groupId>redis.clients</groupId>
            <artifactId>jedis</artifactId>
            <version>5.1.0</version>
        </dependency>
        <dependency>
            <groupId>io.opentelemetry</groupId>
            <artifactId>opentelemetry-api</artifactId>
            <version>1.36.0</version>
        </dependency>
        <dependency>
            <groupId>io.opentelemetry</groupId>
            <artifactId>opentelemetry-sdk</artifactId>
            <version>1.36.0</version>
        </dependency>
        <dependency>
            <groupId>com.google.cloud.opentelemetry</groupId>
            <artifactId>exporter-trace</artifactId>
            <version>0.28.0</version>
        </dependency>
        <dependency>
            <groupId>com.google.cloud.opentelemetry</groupId>
            <artifactId>exporter-metrics</artifactId>
            <version>0.28.0</version>
        </dependency>
    
        <!-- Testing Dependencies -->
        <dependency>
            <groupId>junit</groupId>
            <artifactId>junit</artifactId>
            <version>4.13.2</version>
            <scope>test</scope>
        </dependency>
        <dependency>
            <groupId>org.mockito</groupId>
            <artifactId>mockito-core</artifactId>
            <version>4.11.0</version>
            <scope>test</scope>
        </dependency>
        <dependency>
            <groupId>org.slf4j</groupId>
            <artifactId>slf4j-simple</artifactId>
            <version>1.7.36</version>
            <scope>test</scope>
        </dependency>
    </dependencies>
  2. Erstellen Sie eine RedisTelemetryApp.java-Datei und fügen Sie ihr den folgenden Code hinzu, um die clientseitigen Messwerte zu aktivieren:

    import com.google.cloud.opentelemetry.metric.GoogleCloudMetricExporter;
    import com.google.cloud.opentelemetry.trace.TraceExporter;
    import io.opentelemetry.api.OpenTelemetry;
    import io.opentelemetry.api.common.AttributeKey;
    import io.opentelemetry.api.common.Attributes;
    import io.opentelemetry.api.metrics.DoubleHistogram;
    import io.opentelemetry.api.metrics.LongCounter;
    import io.opentelemetry.api.metrics.Meter;
    import io.opentelemetry.api.trace.Span;
    import io.opentelemetry.api.trace.Tracer;
    import io.opentelemetry.sdk.OpenTelemetrySdk;
    import io.opentelemetry.sdk.metrics.export.MetricExporter;
    import io.opentelemetry.sdk.metrics.SdkMeterProvider;
    import io.opentelemetry.sdk.metrics.export.PeriodicMetricReader;
    import io.opentelemetry.sdk.trace.SdkTracerProvider;
    import io.opentelemetry.sdk.trace.export.BatchSpanProcessor;
    import io.opentelemetry.sdk.trace.export.SpanExporter;
    import redis.clients.jedis.Jedis;
    import redis.clients.jedis.JedisPool;
    import redis.clients.jedis.JedisPoolConfig;
    import redis.clients.jedis.exceptions.JedisConnectionException;
    
    import java.time.Duration;
    import java.util.function.Function;
    
    /**
     * Sample application demonstrating client-side metrics and tracing for
     * Google Cloud Memorystore for Redis.
     */
    public final class RedisTelemetryApp {
        /** Attribute key for Redis operation names. */
        private static final AttributeKey<String> ATTR_OPERATION =
                AttributeKey.stringKey("operation");
    
        /** Maximum number of Redis reconnection attempts. */
        private static final int MAX_RETRIES = 3;
    
        /** Maximum total connections for the Jedis pool. */
        private static final int POOL_MAX_TOTAL = 20;
    
        /** Interval in seconds for exporting metrics to Google Cloud. */
        private static final long METRIC_INTERVAL_SECONDS = 10L;
    
        /** Base multiplier for exponential backoff sleep (in milliseconds). */
        private static final long RETRY_BACKOFF_BASE_MS = 100L;
    
        /** Conversion factor from Nanoseconds to Milliseconds. */
        private static final double NANO_TO_MS = 1_000_000.0;
    
        /** Default Redis port. */
        private static final int DEFAULT_REDIS_PORT = 6379;
    
        /** OpenTelemetry Tracer instance for recording trace spans. */
        private static Tracer tracer;
    
        /** OpenTelemetry Histogram for Redis round-trip time. */
        private static DoubleHistogram rttHist;
    
        /** OpenTelemetry Histogram for pool blocking latency. */
        private static DoubleHistogram clientBlockHist;
    
        /** OpenTelemetry Histogram for application logic blocking latency. */
        private static DoubleHistogram appBlockHist;
    
        /** OpenTelemetry Counter for Redis reconnection retry events. */
        private static LongCounter retryCounter;
    
        /** OpenTelemetry Counter for Redis connectivity errors. */
        private static LongCounter connErrorCounter;
    
        /** Shared Jedis connection pool. */
        private static JedisPool jedisPool;
    
        /**
         * Private constructor to prevent instantiation of this utility class.
         */
        private RedisTelemetryApp() {
        }
    
        /**
         * Main entry point for running the sample application.
         *
         * @param args Command line arguments (not used).
         */
        public static void main(final String[] args) {
            setupTelemetry();
    
            final String host = System.getenv()
                    .getOrDefault("REDISHOST", "localhost");
            final int port = Integer.parseInt(System.getenv()
                    .getOrDefault("REDISPORT",
                            String.valueOf(DEFAULT_REDIS_PORT)));
    
            final JedisPoolConfig poolConfig = new JedisPoolConfig();
            poolConfig.setMaxTotal(POOL_MAX_TOTAL);
            poolConfig.setBlockWhenExhausted(true);
            jedisPool = new JedisPool(poolConfig, host, port);
    
            try {
                run();
            } finally {
                if (jedisPool != null) {
                    jedisPool.close();
                }
            }
        }
    
        /**
         * Executes the core business logic of reading and writing to Redis.
         *
         * @return The string retrieved from the Redis 'get' operation.
         */
        static String run() {
            final Span span = tracer.spanBuilder("process_user_span")
                    .startSpan();
            try {
                smartRedisCall("set_user", jedis ->
                        jedis.set("user:123", "active"));
    
                final String result = smartRedisCall("get_user", jedis ->
                        jedis.get("user:123"));
                System.out.println("Retrieved: " + result);
                return result;
            } catch (Exception e) {
                span.recordException(e);
                throw e;
            } finally {
                span.end();
            }
        }
    
        /**
         * Injects mocked or no-op telemetry and pool instances for unit testing.
         *
         * @param pool                The mocked or test JedisPool instance.
         * @param testOpenTelemetry The OpenTelemetry instance to use for testing.
         */
        static void initForTest(
                final JedisPool pool,
                final OpenTelemetry testOpenTelemetry) {
            jedisPool = pool;
            tracer = testOpenTelemetry.getTracer("jedis.client");
            final Meter meter = testOpenTelemetry.getMeter("jedis.metrics");
    
            rttHist = meter.histogramBuilder("redis_client_rtt")
                    .setUnit("ms").build();
            clientBlockHist = meter
                    .histogramBuilder("redis_client_blocking_latency")
                    .setUnit("ms").build();
            appBlockHist = meter
                    .histogramBuilder("redis_application_blocking_latency")
                    .setUnit("ms").build();
            retryCounter = meter.counterBuilder("redis_retry_count").build();
            connErrorCounter = meter
                    .counterBuilder("redis_connectivity_error_count")
                    .build();
    
            retryCounter.add(0, Attributes.of(ATTR_OPERATION, "startup"));
            connErrorCounter.add(0, Attributes.of(ATTR_OPERATION, "startup"));
        }
    
        /**
         * Configures the production OpenTelemetry SDK to export Traces and Metrics
         * to Google Cloud Operations.
         */
        private static void setupTelemetry() {
            final SpanExporter traceExporter =
                    TraceExporter.createWithDefaultConfiguration();
            final SdkTracerProvider tracerProvider =
                    SdkTracerProvider.builder()
                            .addSpanProcessor(
                                    BatchSpanProcessor.builder(traceExporter)
                                            .build())
                            .build();
    
            final MetricExporter metricExporter =
                    GoogleCloudMetricExporter.createWithDefaultConfiguration();
            final SdkMeterProvider meterProvider =
                    SdkMeterProvider.builder()
                            .registerMetricReader(
                                    PeriodicMetricReader.builder(metricExporter)
                                            .setInterval(Duration.ofSeconds(
                                                    METRIC_INTERVAL_SECONDS))
                                            .build())
                            .build();
    
            final OpenTelemetry openTelemetry = OpenTelemetrySdk.builder()
                    .setTracerProvider(tracerProvider)
                    .setMeterProvider(meterProvider)
                    .buildAndRegisterGlobal();
    
            tracer = openTelemetry.getTracer("jedis.client");
            final Meter meter = openTelemetry.getMeter("jedis.metrics");
    
            rttHist = meter.histogramBuilder("redis_client_rtt")
                    .setUnit("ms").build();
            clientBlockHist = meter
                    .histogramBuilder("redis_client_blocking_latency")
                    .setUnit("ms").build();
            appBlockHist = meter
                    .histogramBuilder("redis_application_blocking_latency")
                    .setUnit("ms").build();
            retryCounter = meter.counterBuilder("redis_retry_count").build();
            connErrorCounter = meter
                    .counterBuilder("redis_connectivity_error_count")
                    .build();
    
            retryCounter.add(0, Attributes.of(ATTR_OPERATION, "startup"));
            connErrorCounter.add(0, Attributes.of(ATTR_OPERATION, "startup"));
        }
    
        /**
         * Wraps a Redis operation with latency metrics, reconnection retry logic,
         * and trace spans.
         *
         * @param <T>           The return type of the Redis operation.
         * @param operationName The name of the operation for metric attributes.
         * @param operation     The Redis command lambda to execute safely.
         * @return The return value from the Redis command.
         */
        private static <T> T smartRedisCall(
                final String operationName,
                final Function<Jedis, T> operation) {
            int attempt = 0;
            final Attributes attrs = Attributes.of(ATTR_OPERATION,
                    operationName);
    
            final Span span = tracer.spanBuilder(operationName).startSpan();
    
            try {
                while (attempt < MAX_RETRIES) {
                    final long poolStart = System.nanoTime();
                    try (Jedis jedis = jedisPool.getResource()) {
                        clientBlockHist.record((System.nanoTime() - poolStart)
                                / NANO_TO_MS, attrs);
    
                        final long reqStart = System.nanoTime();
                        final T response = operation.apply(jedis);
                        rttHist.record((System.nanoTime() - reqStart)
                                / NANO_TO_MS, attrs);
    
                        final long appStart = System.nanoTime();
                        @SuppressWarnings("unused")
                        final String dummy = String.valueOf(response);
                        appBlockHist.record((System.nanoTime() - appStart)
                                / NANO_TO_MS, attrs);
    
                        return response;
                    } catch (JedisConnectionException e) {
                        attempt++;
                        connErrorCounter.add(1, attrs);
                        retryCounter.add(1, attrs);
                        span.recordException(e);
                        if (attempt >= MAX_RETRIES) {
                            throw e;
                        }
                        try {
                            Thread.sleep((long) (Math.pow(2, attempt)
                                    * RETRY_BACKOFF_BASE_MS));
                        } catch (InterruptedException ie) {
                            Thread.currentThread().interrupt();
                        }
                    }
                }
                return null;
            } finally {
                span.end();
            }
        }
    }
  3. Führen Sie Ihre Anwendung mindestens eine Minute lang aus, damit der Exporter genügend Zeit hat, die veröffentlichten Messwerte zu verarbeiten und an Monitoring zu senden.

Node.js

  1. Führen Sie die folgenden Befehle in Ihrem Terminal aus, um die erforderlichen OpenTelemetry- und Google Cloud Exporter Abhängigkeiten zu installieren:

      npm install redis@^4.6.0 @opentelemetry/api@^1.9.0
      @opentelemetry/sdk-trace-node@^2.1.0
      @opentelemetry/sdk-trace-base@^2.1.0
      @opentelemetry/sdk-metrics@^2.1.0
      @opentelemetry/instrumentation@^0.205.0
      @opentelemetry/instrumentation-redis@^0.67.0
      @google-cloud/opentelemetry-cloud-trace-exporter@^3.0.0
      @google-cloud/opentelemetry-cloud-monitoring-exporter@^0.21.0
      @opentelemetry/resources@^2.1.0
  2. Erstellen Sie eine server.js-Datei und fügen Sie ihr den folgenden Code hinzu, um die clientseitigen Messwerte zu aktivieren:

    
    'use strict';
    
    const {trace, metrics} = require('@opentelemetry/api');
    const {NodeTracerProvider} = require('@opentelemetry/sdk-trace-node');
    const {BatchSpanProcessor} = require('@opentelemetry/sdk-trace-base');
    const {
      TraceExporter,
    } = require('@google-cloud/opentelemetry-cloud-trace-exporter');
    const {
      MeterProvider,
      PeriodicExportingMetricReader,
    } = require('@opentelemetry/sdk-metrics');
    const {
      MetricExporter,
    } = require('@google-cloud/opentelemetry-cloud-monitoring-exporter');
    const {RedisInstrumentation} = require('@opentelemetry/instrumentation-redis');
    const {registerInstrumentations} = require('@opentelemetry/instrumentation');
    const {performance} = require('perf_hooks');
    
    // FIX: Pass spanProcessors in the constructor options for NodeTracerProvider in SDK 2.x
    const provider = new NodeTracerProvider({
      spanProcessors: [new BatchSpanProcessor(new TraceExporter())],
    });
    provider.register();
    
    registerInstrumentations({
      instrumentations: [new RedisInstrumentation()],
    });
    
    const redis = require('redis');
    
    const metricExporter = new MetricExporter();
    const metricReader = new PeriodicExportingMetricReader({
      exporter: metricExporter,
      exportIntervalMillis: 10000,
    });
    const meterProvider = new MeterProvider({readers: [metricReader]});
    metrics.setGlobalMeterProvider(meterProvider);
    
    const tracer = trace.getTracer('redis.client.node');
    const meter = metrics.getMeter('redis.metrics.node');
    
    const rttHist = meter.createHistogram('redis_client_rtt', {unit: 'ms'});
    const appBlockHist = meter.createHistogram(
      'redis_application_blocking_latency',
      {unit: 'ms'}
    );
    const retryCounter = meter.createCounter('redis_retry_count');
    const connErrorCounter = meter.createCounter('redis_connectivity_error_count');
    
    retryCounter.add(0, {operation: 'startup'});
    connErrorCounter.add(0, {operation: 'startup'});
    
    const REDISHOST = process.env.REDISHOST || 'localhost';
    const REDISPORT = process.env.REDISPORT || 6379;
    
    const client = redis.createClient({
      socket: {
        host: REDISHOST,
        port: REDISPORT,
        reconnectStrategy: retries => {
          connErrorCounter.add(1, {error: 'socket_reconnect'});
          if (retries > 5) return new Error('Max retries reached');
          return Math.min(retries * 100, 3000);
        },
      },
    });
    client.on('error', err => console.log('Redis Client Error', err));
    
    async function smartRedisCall(operationName, func, ...args) {
      let attempt = 0;
      while (attempt < 3) {
        try {
          const reqStart = performance.now();
          const response = await func(...args);
          rttHist.record(performance.now() - reqStart, {operation: operationName});
    
          const appParseStart = performance.now();
          // eslint-disable-next-line no-unused-vars
          const _ = String(response);
          appBlockHist.record(performance.now() - appParseStart, {
            operation: operationName,
          });
    
          return response;
        } catch (e) {
          attempt++;
          retryCounter.add(1, {operation: operationName});
          if (attempt >= 3) throw e;
          await new Promise(resolve =>
            setTimeout(resolve, Math.pow(2, attempt) * 100)
          );
        }
      }
    }
    
    async function main() {
      await client.connect();
    
      await tracer.startActiveSpan('process_user_span', async span => {
        try {
          // Simple write and read operations
          await smartRedisCall(
            'set_user',
            client.set.bind(client),
            'user:123',
            'active'
          );
    
          const result = await smartRedisCall(
            'get_user',
            client.get.bind(client),
            'user:123'
          );
          console.log('Retrieved:', result);
        } catch (e) {
          span.recordException(e);
        } finally {
          span.end();
        }
      });
    
      await client.quit();
      await provider.forceFlush();
      await meterProvider.forceFlush();
    }
    
    // Only run the script automatically if it is executed directly (e.g. `node server.js`)
    if (require.main === module) {
      main().catch(console.error);
    }
    
    // Export for testability
    module.exports = {
      main,
      smartRedisCall,
    };
    
  3. Führen Sie Ihre Anwendung mindestens eine Minute lang aus, damit der Exporter genügend Zeit hat, die veröffentlichten Messwerte zu verarbeiten und an Monitoring zu senden.

Python

  1. Führen Sie die folgenden Befehle in Ihrem Terminal aus, um die erforderlichen OpenTelemetry- und Google Cloud Exporter Abhängigkeiten zu installieren:

      pip install redis==7.0.1 opentelemetry-api==1.39.1
      opentelemetry-sdk==1.39.1
      opentelemetry-instrumentation-redis==0.60b1
      opentelemetry-exporter-gcp-trace==1.11.0
      opentelemetry-exporter-gcp-monitoring==1.11.0a0
  2. Erstellen Sie eine main.py-Datei und fügen Sie ihr den folgenden Code hinzu, um die clientseitigen Messwerte zu aktivieren:

    import os
    import time
    
    from opentelemetry import metrics, trace
    from opentelemetry.exporter.cloud_monitoring import (
        CloudMonitoringMetricsExporter,
    )
    from opentelemetry.exporter.cloud_trace import CloudTraceSpanExporter
    from opentelemetry.instrumentation.redis import RedisInstrumentor
    from opentelemetry.sdk.metrics import MeterProvider
    from opentelemetry.sdk.metrics.export import PeriodicExportingMetricReader
    from opentelemetry.sdk.trace import TracerProvider
    from opentelemetry.sdk.trace.export import BatchSpanProcessor
    import redis
    from redis.exceptions import ConnectionError, TimeoutError
    
    
    
    
    def init_telemetry():
        """Initializes OpenTelemetry with GCP Exporters and returns the SDK objects."""
        # 1. Initialize Tracing
        tracer_provider = TracerProvider()
        tracer_provider.add_span_processor(
            BatchSpanProcessor(CloudTraceSpanExporter())
        )
        trace.set_tracer_provider(tracer_provider)
        tracer = trace.get_tracer("redis.client")
    
        # 2. Initialize Metrics
        metrics_exporter = CloudMonitoringMetricsExporter()
        metric_reader = PeriodicExportingMetricReader(
            metrics_exporter, export_interval_millis=10000
        )
        meter_provider = MeterProvider(metric_readers=[metric_reader])
        metrics.set_meter_provider(meter_provider)
        meter = metrics.get_meter("redis.metrics")
    
        # Bundle all metric handlers safely into a dictionary
        redis_metrics = {
            "rtt_hist": meter.create_histogram("redis_client_rtt", unit="ms"),
            "client_block_hist": meter.create_histogram(
                "redis_client_blocking_latency", unit="ms"
            ),
            "app_block_hist": meter.create_histogram(
                "redis_application_blocking_latency", unit="ms"
            ),
            "retry_counter": meter.create_counter("redis_retry_count"),
            "conn_error_counter": meter.create_counter(
                "redis_connectivity_error_count"
            ),
        }
    
        redis_metrics["retry_counter"].add(0, {"operation": "startup"})
        redis_metrics["conn_error_counter"].add(0, {"operation": "startup"})
    
        # 3. Setup Redis Auto-Instrumentation
        RedisInstrumentor().instrument()
    
        return tracer, redis_metrics, tracer_provider, meter_provider
    
    
    def init_redis_pool():
        """Initializes and returns the Redis ConnectionPool and Client."""
        redis_host = os.environ.get("REDISHOST", "localhost")
        redis_port = int(os.environ.get("REDISPORT", 6379))
    
        redis_pool = redis.ConnectionPool(
            host=redis_host,
            port=redis_port,
            max_connections=10,
            decode_responses=True,
        )
        redis_client = redis.Redis(connection_pool=redis_pool)
        return redis_pool, redis_client
    
    
    def smart_redis_call(
        operation_name, func, redis_pool, metrics, *args, **kwargs
    ):
        """Executes a Redis operation with metrics and retry handling (No Globals!)."""
        max_retries = 3
        attempt = 0
    
        pool_start = time.time()
        try:
            conn = redis_pool.get_connection()
            redis_pool.release(conn)
        except Exception:
            pass
    
        if metrics and metrics.get("client_block_hist"):
            metrics["client_block_hist"].record(
                (time.time() - pool_start) * 1000, {"operation": operation_name}
            )
    
        while attempt < max_retries:
            try:
                req_start = time.time()
                response = func(*args, **kwargs)
    
                if metrics and metrics.get("rtt_hist"):
                    metrics["rtt_hist"].record(
                        (time.time() - req_start) * 1000,
                        {"operation": operation_name},
                    )
    
                app_start = time.time()
                _ = str(response)
    
                if metrics and metrics.get("app_block_hist"):
                    metrics["app_block_hist"].record(
                        (time.time() - app_start) * 1000,
                        {"operation": operation_name},
                    )
    
                return response
    
            except (ConnectionError, TimeoutError) as e:
                attempt += 1
                if metrics and metrics.get("conn_error_counter"):
                    metrics["conn_error_counter"].add(
                        1, {"operation": operation_name}
                    )
                if metrics and metrics.get("retry_counter"):
                    metrics["retry_counter"].add(1, {"operation": operation_name})
                if attempt >= max_retries:
                    raise e
                time.sleep((2**attempt) * 0.1)
    
    if __name__ == "__main__":
        tracer, redis_metrics, tracer_provider, meter_provider = init_telemetry()
        redis_pool, redis_client = init_redis_pool()
    
        if tracer:
            with tracer.start_as_current_span("process_user_span"):
                try:
                    # Simple write and read operations
                    smart_redis_call(
                        "set_user",
                        redis_client.set,
                        redis_pool,
                        redis_metrics,
                        "user:123",
                        "active",
                    )
    
                    result = smart_redis_call(
                        "get_user",
                        redis_client.get,
                        redis_pool,
                        redis_metrics,
                        "user:123",
                    )
                    print(f"Retrieved: {result}")
                except Exception as e:
                    print(f"Error: {e}")
    
            tracer_provider.force_flush()
            meter_provider.force_flush()
  3. Führen Sie Ihre Anwendung mindestens eine Minute lang aus, damit der Exporter genügend Zeit hat, die veröffentlichten Messwerte zu verarbeiten und an Monitoring zu senden.

Messwerte in Monitoring ansehen

Nachdem Sie clientseitige Messwerte aktiviert und Ihre Anwendung mindestens eine Minute lang ausgeführt haben, damit der Exporter genügend Zeit hat, Messwerte zu verarbeiten und an Monitoring zu senden, können Sie Ihre Messwerte in Monitoring visualisieren, nach Vorgang oder Instanz gruppieren und Aggregatoren anwenden, um die Leistung Ihrer Anwendung zu überwachen.

So rufen Sie Messwerte in Monitoring auf:

  1. Wechseln Sie in der Google Cloud Console zur Seite Metrics Explorer.

    Zum Metrics Explorer

  2. Wählen Sie Ihr Google Cloud Projekt aus.

  3. Klicken Sie auf Messwert auswählen.

  4. Suchen Sie nach workload.googleapis.com/redis.

  5. Wählen Sie einen clientseitigen Messwert aus. Gruppieren Sie die Daten nach Bedarf nach operation und instance und wählen Sie einen Aggregator aus. Weitere Optionen finden Sie unter Messwerte bei Verwendung von Metrics Explorer auswählen.

Verteilte Traces in Trace ansehen

Nachdem Ihre Anwendung mit dem Exportieren von Daten begonnen hat, können Sie mit Trace den gesamten Anfrage-Antwort-Zyklus Ihrer Redis-Befehle visualisieren. Wenn Sie Ihre verteilten Traces in Trace ansehen, können Sie Engpässe diagnostizieren und so die genaue Quelle der Latenz in Ihrer Anwendung schnell isolieren.

So rufen Sie verteilte Traces in Trace auf:

  1. Rufen Sie in der Google Cloud Console die Seite Trace Explorer auf.

    Trace Explorer aufrufen

  2. Wählen Sie einen aktuellen Trace aus, der durch einen Punkt im Streudiagramm dargestellt wird.

  3. Untersuchen Sie die Wasserfallansicht, um die Quelle der Latenz zu isolieren, indem Sie die folgenden Engpässe ermitteln:

    • Gesamtdauer der Anfrage: Der Balken der obersten Ebene (übergeordnet) zeigt die Gesamt zeit an, die Sie warten müssen, bis der Vorgang abgeschlossen ist.

    • Netzwerk- und Serverlatenz (Round-Trip-Zeit, RTT): Die untergeordneten Balken (z. B. mit der Bezeichnung GET oder SET) zeigen die Zeit an, die der Befehl für die Übertragung im Netzwerk und die Ausführung auf dem Memorystore for Redis-Server benötigt hat.

    • Blockierung der Clientverbindung: Wenn vor dem Beginn des untergeordneten Redis-Spans eine große, leere horizontale Lücke vorhanden ist, wartet der Anwendungsthread auf eine verfügbare TCP-Verbindung aus dem Verbindungspool.

    • Blockierung der Anwendungsparser: Wenn nach dem Ende des untergeordneten Redis-Spans eine große, leere horizontale Lücke vorhanden ist, hat die Anwendung Schwierigkeiten, die zurückgegebene Nutzlast zu parsen oder zu verarbeiten. Dies tritt häufig bei JSON-Strings mit mehreren Megabyte auf.

    • Wiederholungen: Wenn Sie mehrere kurze untergeordnete Spans für denselben Befehl sehen, die nacheinander im selben übergeordneten Trace auftreten, kann es sein, dass bei Ihrem Client Netzwerkpaketverluste auftreten und er seine exponentielle Backoff-Wiederholungsschleife auslösen muss.

Fehlerbehebung

In diesem Abschnitt werden häufige Leistungsprobleme aufgeführt, die Sie mithilfe clientseitiger Messwerte ermitteln können. Außerdem werden die Ursachen dieser Probleme erläutert und Anleitungen zur Fehlerbehebung gegeben.

Problem Ursache Fehlerbehebung

In Ihrer Anwendung tritt ein plötzlicher Latenzspike auf, aber Memorystore for Redis scheint völlig fehlerfrei zu sein.

  • workload.googleapis.com/
    redis_client_blocking_latency
    (clientseitiger Messwert): Spike
  • workload.googleapis.com/redis_client_rtt (clientseitiger Messwert): niedrig / typisch
  • redis.googleapis.com/commands/
    usec_per_call
    (Memorystore for Redis-Servermesswert): niedrig / typisch
  • redis.googleapis.com/clients/connected (Memorystore for Redis-Servermesswert): bleibt bei einer bestimmten Zahl
Der Engpass befindet sich ausschließlich in Ihrer Anwendung. Ihre Threads versuchen, Redis-Befehle auszuführen, aber der Verbindungspool ist vollständig ausgeschöpft. Die hohe redis_client_blocking_latency gibt die Zeit an, die Ihr Code auf einen verfügbaren TCP-Socket wartet, bevor der Befehl an das Netzwerk gesendet wird. Erhöhen Sie die Grenzwerte für die Größe des Verbindungspools in Ihrer Redis-Clientkonfiguration, um den höheren gleichzeitigen Traffic zu bewältigen (z. B. MaxActive für Go, MaxTotal für Java oder max_connections für Node.js und Python).

Die Anfrage wird abgeschlossen, aber der Endpunkt benötigt deutlich länger als erwartet. Es gibt keine Probleme mit dem Zustand Ihres Netzwerks oder Servers.

  • workload.googleapis.com/
    redis_application_blocking_latency
    (clientseitiger Messwert): Spike
  • workload.googleapis.com/redis_client_rtt (clientseitig metric): niedrig / typisch
  • redis.googleapis.com/commands/
    usec_per_call
    (Memorystore for Redis-Servermesswert): niedrig / typisch
  • redis.googleapis.com/stats/
    network_traffic
    (Ausgehende Bytes) (Memorystore for Redis-Servermesswert): starke Spikes
Memorystore for Redis führt den Befehl aus und das Netzwerk überträgt die Nutzlast schnell (niedrige RTT). Die zurückgegebene Nutzlast ist jedoch groß (z. B. ein 15 MB großer JSON-String). In Ihrer Anwendung tritt eine hohe redis_application_blocking_latency auf, da die Anwendung übermäßig viele Ressourcen verbraucht, während sie Arbeitsspeicher zuweist und diesen großen String in ein Objekt deserialisiert. Optimieren Sie Ihr Datenmodell. Speichern Sie keine riesigen JSON-Blobs in einzelnen Schlüsseln. Unterteilen Sie die Daten mit Redis-Hashes (HSET) und rufen Sie mit HGET oder HMGET nur die benötigten Felder ab.

Die Latenz Ihrer nutzerorientierten Anwendung steigt, aber Ihre Redis-Messwerte melden eine niedrige Serverlatenz und typische Auscheckvorgänge im Verbindungspool.

  • workload.googleapis.com/redis_retry_count (clientseitiger Messwert): Spike
  • workload.googleapis.com/
    redis_connectivity_error_count
    (clientseitiger Messwert): kann vorübergehende Erhöhungen aufweisen
  • workload.googleapis.com/redis_client_rtt (clientseitiger Messwert): niedrig / typisch für erfolgreiche Anfragen
  • redis.googleapis.com/commands/
    usec_per_call
    (Memorystore for Redis-Servermesswert): niedrig / typisch
Da redis_client_rtt nur die RTT erfolgreicher Anfragen erfasst, wird die Zeitüberschreitungsdauer eines fehlgeschlagenen Pakets nicht berücksichtigt. Wenn in Ihrer Anwendung vorübergehende Paketverluste oder TCP Zurücksetzungen auftreten, erhöht die Wiederholungslogik Ihres instrumentierten Clients redis_retry_count und löst die exponentielle Backoff-Schleife aus. Dadurch wird eine Wartezeit zwischen den Versuchen eingeführt (z. B. 100ms, 200ms, oder 400ms). Der Nutzer erlebt eine hohe Gesamtlatenz, aber die zugrunde liegende Ursache ist ein Netzwerk paketverlust, der clientseitige Wartezeiten auslöst. Prüfen Sie Ihre VPC-Flusslogs auf verworfene Pakete, Bandbreitenbegrenzung, oder regionsübergreifende Routinganomalien. Wenn aggressive Zeitüberschreitungen auftreten, müssen die Zeitüberschreitungen für die Clientverbindung (socket_timeout oder connect_timeout) größer als die erwartete RTT sein, um vorübergehende Netzwerklatenz zu berücksichtigen.

Alles stoppt und alle Ebenen der Telemetriepipeline melden hohe Latenz.

  • workload.googleapis.com/redis_client_rtt (clientseitig metric): high
  • redis.googleapis.com/commands/
    usec_per_call
    (Memorystore for Redis-Servermesswert): hoch
  • redis.googleapis.com/stats/
    cpu_utilization_main_thread
    (Memorystore for Redis-Servermesswert): hoch (z. B. nahe 1 s/s, oder 100%)
  • Trace-Wasserfallansicht: zeigt einen Befehl, der viel Zeit in Anspruch nimmt
Redis ist single-threaded. Wenn Sie einen Befehl mit der Zeitkomplexität O(N) ausführen, z. B. KEYS *, SMEMBERS für eine große Menge oder HGETALL für einen Hash mit Millionen von Feldern, wird die Redis-Engine angehalten, um diese Anfrage zu bearbeiten. Während dieser Befehl ausgeführt wird, werden alle anderen Anwendungsanfragen in die Warteschlange gestellt, was zu einem systemweiten Latenz spike führt. Da Ihre benutzerdefinierte redis_client_rtt mit der Serverlatenz (commands/usec_per_call) übereinstimmt, ist der Server, auf dem der Befehl ausgeführt wird, der Engpass.

Öffnen Sie Trace und sehen Sie sich die Redis-Befehle in den langsamen Spans an, um zu ermitteln, welche Abfrage die Blockierung verursacht. Ersetzen Sie blockierende Befehle durch nicht blockierende in Ihrem Code.

Wenn Sie große Datensätze inkrementell durchlaufen möchten, ohne den Serverthread zu sperren, verwenden Sie SCAN, SSCAN, oder HSCAN.

Ihre Anwendung meldet eine gleichbleibend erhöhte Baseline-Latenz für alle Redis-Befehle, auch wenn der Traffic gering ist.

  • workload.googleapis.com/redis_client_rtt (clientseitiger Messwert): konstant erhöht (p50 und p99 sind beide ~30–100 ms+)
  • redis.googleapis.com/commands/
    usec_per_call
    (Memorystore for Redis-Servermesswert): extrem niedrig (< 1 ms)
  • workload.googleapis.com/
    redis_client_blocking_latency

    und redis_application_blocking_latency: niedrig / typisch
Der Redis-Server führt Befehle sofort aus, aber Ihre Anwendung und Ihre Instanz werden in verschiedenen Regionen bereitgestellt (z. B. us-central1 und us-east1). Jedes Netzwerkpaket muss zwischen diesen geografischen Rechenzentren über die physische Google Cloud-Infrastruktur übertragen werden. Dies führt zu einer obligatorischen regionsübergreifenden Latenzstrafe in Höhe der Lichtgeschwindigkeit für jede Round-Trip-Zeit. Wenn Sie die Latenz reduzieren möchten, stellen Sie Ihre Anwendung in derselben Region und Zone wie Ihre Instanz bereit. In der Google Cloud Console können Sie die Region Ihrer Anwendung und Instanz ansehen.

Nächste Schritte