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

‫Memorystore for Redis מספק מדדים בזמן אמת בצד השרת למעקב אחרי קצב העברת הנתונים, ניצול המעבד ושימוש בזיכרון, אבל יכול להיות שהנתונים האלה לבדם לא יסבירו למה אפליקציית הלקוח שלכם חווה חביון גבוה במערכות מורכבות ומבוזרות.

מדדים בצד הלקוח פותרים את הבעיה הזו כי הם מספקים שקיפות לגבי המחזור המלא של בקשות ותגובות. הם מודדים פקודה מהרגע שבו האפליקציה מתחילה אותה ועד שהאפליקציה מעבדת את התגובה. בעזרת הנתונים האלה אפשר לקבוע בצורה מדויקת אם זמן האחזור נובע מהלוגיקה של האפליקציה, מנתיב הרשת או משרת Redis.

לפני שמתחילים

מוודאים שאפליקציית הלקוח משתמשת בחשבון שירות, ושמוקצים לה תפקידי ניהול הזהויות והרשאות הגישה (IAM) הבאים:

  • roles/cloudtrace.agent (Cloud Trace Agent)
  • roles/monitoring.metricWriter (Monitoring Metric Writer)

מידע נוסף על מתן תפקידים זמין במדריך למתחילים איך נותנים תפקידים ב-IAM באמצעות מסוף Google Cloud .

הפעלת Cloud Monitoring API

כדי לייצא מדדים בצד הלקוח אל Monitoring, צריך להפעיל את Monitoring API באפליקציה. ייצוא של המדדים האלה ל-Monitoring והצגה שלהם מאפשרים לכם לזהות את שורש הבעיה של צווארי בקבוק כדי לקבוע מהיכן נובעת ההשהיה.

כדי להפעיל את Monitoring API:

  1. במסוף Google Cloud , נכנסים לדף APIs & Services.

    עוברים אל 'ממשקי API ושירותים'

  2. בוחרים את הפרויקט שבו יצרתם את מכונת Memorystore for Redis.

  3. לוחצים על Enable APIs and services.

  4. חיפוש של monitoring.

  5. בתוצאות החיפוש, לוחצים על Cloud Monitoring API.

  6. אם מופיע הכיתוב API enabled, סימן שממשק ה-API כבר מופעל. אחרת, לוחצים על הפעלה.

הפעלת Cloud Trace API

כדי לראות עקבות מבוזרים ב-Trace, צריך להפעיל את Trace API. אחר כך תוכלו להשתמש ב-Trace Explorer כדי לראות את העקבות האלה, לאבחן צווארי בקבוק ולבודד את מקור זמן האחזור באפליקציה.

כדי להפעיל את Trace API:

  1. במסוף Google Cloud , נכנסים לדף APIs & Services.

    עוברים אל 'ממשקי API ושירותים'

  2. בוחרים את הפרויקט שבו יצרתם את מכונת Memorystore for Redis.

  3. לוחצים על Enable APIs and services.

  4. חיפוש של trace.

  5. בתוצאות החיפוש, לוחצים על Cloud Trace API.

  6. אם מופיע הכיתוב API enabled, סימן שממשק ה-API כבר מופעל. אחרת, לוחצים על הפעלה.

הפעלה של מדדים בצד הלקוח

כדי להפעיל מדדים בצד הלקוח, מוסיפים את OpenTelemetry SDK, את כלי הייצוא של Cloud Monitoring ואת כלי הייצוא של Cloud Trace לקוד של האפליקציה. המדדים נאספים על ידי OpenTelemetry, שפועל ישירות בספריית הלקוח של Redis באפליקציה. כך האפליקציה יכולה לתעד נקודות נתונים של זמן האחזור ולייצא אותן ל-Monitoring ול-Trace לצורך ויזואליזציה.

כדי להפעיל מדדים בצד הלקוח, אפשר להשתמש ב-Go,‏ Java,‏ Node.js או Python. בכרטיסיות הבאות מופיע מידע על הפעלת המדדים לכל שפה.

המשך

  1. כדי להתקין את יחסי התלות הנדרשים של OpenTelemetry ושל Google Cloud כלי לייצוא נתונים מריצים את הפקודות הבאות בטרמינל:

      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. כדי להפעיל את המדדים בצד הלקוח, יוצרים קובץ main.go ומוסיפים לו את הקוד הבא:

    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. מריצים את האפליקציה למשך דקה לפחות כדי לתת לכלי הייצוא מספיק זמן לאגד את המדדים שפורסמו ולשלוח אותם ל-Monitoring.

Java

  1. כדי להתקין את יחסי התלות הנדרשים של OpenTelemetry ושל Google Cloud exporter מוסיפים את הקוד הבא לקובץ pom.xml של האפליקציה:

    <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. כדי להפעיל את המדדים בצד הלקוח, יוצרים קובץ RedisTelemetryApp.java ומוסיפים לו את הקוד הבא:

    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. מריצים את האפליקציה למשך דקה לפחות כדי לתת לכלי הייצוא מספיק זמן לאגד את המדדים שפורסמו ולשלוח אותם ל-Monitoring.

Node.js

  1. כדי להתקין את יחסי התלות הנדרשים של OpenTelemetry ושל Google Cloud כלי לייצוא נתונים מריצים את הפקודות הבאות בטרמינל:

      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. כדי להפעיל את המדדים בצד הלקוח, יוצרים קובץ server.js ומוסיפים לו את הקוד הבא:

    
    '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. מריצים את האפליקציה למשך דקה לפחות כדי לתת לכלי הייצוא מספיק זמן לאגד את המדדים שפורסמו ולשלוח אותם ל-Monitoring.

Python

  1. כדי להתקין את יחסי התלות הנדרשים של OpenTelemetry ושל Google Cloud כלי לייצוא נתונים מריצים את הפקודות הבאות בטרמינל:

      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. כדי להפעיל את המדדים בצד הלקוח, יוצרים קובץ main.py ומוסיפים את הקוד הבא לאפליקציה:

    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. מריצים את האפליקציה למשך דקה לפחות כדי לתת לכלי הייצוא מספיק זמן לאגד את המדדים שפורסמו ולשלוח אותם ל-Monitoring.

הצגת מדדים ב'מעקב'

אחרי שמפעילים את המדדים בצד הלקוח ומריצים את האפליקציה למשך דקה לפחות כדי לתת לכלי הייצוא מספיק זמן לאגד את המדדים ולשלוח אותם ל-Monitoring, אפשר להשתמש ב-Monitoring כדי להציג את המדדים באופן חזותי, לקבץ אותם לפי פעולה או מופע ולהחיל פונקציות צבירה כדי לעקוב אחרי הביצועים של האפליקציה.

כדי לראות את המדדים ב-Monitoring:

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

    לדף Metrics Explorer

  2. בוחרים את הפרויקט Google Cloud .

  3. לוחצים על בחירת מדד.

  4. חיפוש של workload.googleapis.com/redis.

  5. בוחרים מדד מצד הלקוח. מקבצים את הנתונים לפי operation ו-instance לפי הצורך, ובוחרים פונקציית צבירה. כדי לראות אפשרויות נוספות, אפשר לעיין במאמר איך בוחרים מדדים כשמשתמשים ב-Metrics Explorer.

צפייה בנתוני מעקב מבוזרים ב-Trace

אחרי שהאפליקציה מתחילה לייצא נתונים, אפשר להשתמש ב-Trace כדי להציג באופן חזותי את מחזור הבקשה-תגובה המלא של פקודות Redis. הצגת העקבות המבוזרים ב-Trace מאפשרת לאבחן צווארי בקבוק, כך שתוכלו לבודד במהירות את המקור המדויק של זמן האחזור באפליקציה.

כדי לראות את העקבות המבוזרים ב-Trace:

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

    כניסה לדף Trace Explorer

  2. בוחרים מעקב אחרון שמיוצג על ידי נקודה בתרשים הפיזור.

  3. בודקים את התצוגה של Waterfall כדי לזהות את מקור זמן האחזור. לשם כך, צריך לזהות את צווארי הבקבוק הבאים:

    • משך הבקשה הכולל: בסרגל ברמה העליונה (הורה) מוצג הזמן הכולל שצריך לחכות עד שהפעולה תסתיים.

    • זמן האחזור (RTT) של הרשת והשרת: הפסים המשניים (כמו אלה שמסומנים בתווית GET או SET) מראים את הזמן שחלף מהרגע שהפקודה נשלחה ברשת ועד שהיא הופעלה בשרת Memorystore for Redis.

    • חסימת חיבור לקוח: אם יש רווח אופקי גדול וריק לפני תחילת הטווח המשני של Redis, סימן שהשרשור של האפליקציה תקוע בהמתנה לחיבור TCP זמין ממאגר החיבורים.

    • חסימה של ניתוח האפליקציה: אם יש רווח אופקי גדול וריק אחרי שטווח הצאצא של Redis מסתיים, האפליקציה מתקשה לנתח או לעבד את מטען הייעודי (payload) שמוחזר. זה קורה בדרך כלל עם מחרוזות JSON של כמה מגה-בייט.

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

פתרון בעיות

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

שגיאה מטרה פתרון בעיות

באפליקציה חווים עלייה פתאומית בחביון, אבל נראה ש-Memorystore for Redis תקין לחלוטין.

  • workload.googleapis.com/
    redis_client_blocking_latency
    (מדד מצד הלקוח): עלייה חדה
  • workload.googleapis.com/redis_client_rtt (מדד בצד הלקוח): נמוך / רגיל
  • redis.googleapis.com/commands/
    usec_per_call
    (מדד של שרת Memorystore for Redis): נמוך / אופייני
  • redis.googleapis.com/clients/connected (מדד של שרת Memorystore for Redis): קו ישר במספר מסוים
צוואר הבקבוק נמצא בתוך האפליקציה. השרשורים מנסים להריץ פקודות Redis, אבל מאגר החיבורים מלא. הערך הגבוה redis_client_blocking_latency מייצג את הזמן שבו הקוד ממתין לשקע TCP זמין לפני שהפקודה נשלחת לרשת. כדי לטפל בתנועה המקבילה הגבוהה יותר, צריך להגדיל את מגבלות הגודל של מאגר החיבורים בהגדרות של לקוח Redis (לדוגמה, MaxActive עבור Go,‏ MaxTotal עבור Java או max_connections עבור Node.js ו-Python).

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

  • workload.googleapis.com/
    redis_application_blocking_latency
    (מדד מצד הלקוח): עלייה חדה
  • workload.googleapis.com/redis_client_rtt (מדד בצד הלקוח): נמוך / רגיל
  • redis.googleapis.com/commands/
    usec_per_call
    (מדד של שרת Memorystore for Redis): נמוך / אופייני
  • redis.googleapis.com/stats/
    network_traffic
    ‫(Bytes out) (מדד של שרת Memorystore for Redis): עליות חדות
הפקודה מורצת ב-Memorystore for Redis והמטען הייעודי (payload) מועבר ברשת במהירות (זמן הלוך ושוב נמוך). עם זאת, המטען הייעודי (payload) שמוחזר גדול (לדוגמה, מחרוזת JSON בגודל 15MB). האפליקציה שלך חווה שימוש גבוה בזיכרון redis_application_blocking_latency כי היא צורכת משאבים מוגזמים בזמן הקצאת הזיכרון וביטול הסדרות של המחרוזת הגדולה הזו לאובייקט. אופטימיזציה של מודל הנתונים. לא כדאי לאחסן בלובים גדולים של JSON במפתחות יחידים. מפרקים את הנתונים באמצעות פונקציות hash של Redis ‏ (HSET) ומשתמשים בפקודות HGET או HMGET כדי לאחזר רק את השדות הספציפיים שאתם צריכים.

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

  • workload.googleapis.com/redis_retry_count (מדד מצד הלקוח): עליות חדות
  • workload.googleapis.com/
    redis_connectivity_error_count
    (מדד בצד הלקוח): יכול להיות שיוצגו בו עליות זמניות
  • workload.googleapis.com/redis_client_rtt (מדד בצד הלקוח): נמוך / אופייני לבקשות שהושלמו בהצלחה
  • redis.googleapis.com/commands/
    usec_per_call
    (מדד של שרת Memorystore for Redis): נמוך / אופייני
מכיוון ש-redis_client_rtt מתעד רק את זמן ה-RTT של בקשות מוצלחות, הוא לא משקף את משך הזמן הקצוב לתפוגה של מנות שנכשלו. אם האפליקציה חווה נפילות חבילות זמניות או איפוסים של TCP, לוגיקת הניסיון החוזר של הלקוח שמוגדר בו מעקב מגדילה את הערך של redis_retry_count ומפעילה את לולאת ההשהיה המעריכית לפני ניסיון חוזר שלו. השיטה הזו מוסיפה זמן השהיה בין הניסיונות (למשל, 100ms, 200ms או 400ms). המשתמש חווה זמן השהיה כולל גבוה, אבל שורש הבעיה הוא אובדן מנות ברשת, שגורם להשהיות בצד הלקוח. בודקים ביומני הזרימה של ה-VPC אם יש מנות שהושמטו, הגבלת רוחב פס או חריגות בניתוב בין אזורים. אם אתם חווים התנתקויות רשת לאחר שתם הזמן הקצוב לתפוגה, צריך לוודא שהזמן הקצוב לתפוגה של חיבור הלקוח (socket_timeout או connect_timeout) גדול יותר מה-RTT הצפוי, כדי להתחשב בתנודות רגעיות ברשת.

הכול נעצר וכל השכבות של צינור הטלמטריה מדווחות על זמן טעינה ארוך.

  • workload.googleapis.com/redis_client_rtt (מדד מצד הלקוח): גבוה
  • redis.googleapis.com/commands/
    usec_per_call
    (מדד של שרת Memorystore for Redis): גבוה
  • redis.googleapis.com/stats/
    cpu_utilization_main_thread
    (מדד של שרת Memorystore for Redis): גבוה (לדוגמה, קרוב ל-1 s/s, או 100%)
  • תרשים מפל של מעקב: מציג פקודה שלוקחת הרבה זמן
‫Redis הוא single-threaded. כשמריצים פקודה של מורכבות זמן O(N) – כמו KEYS *, ‏ SMEMBERS על מערך גדול, או HGETALL על hash עם מיליוני שדות – מנוע Redis מושהה כדי למלא את הבקשה. בזמן שהפקודה הזו פועלת, כל בקשה אחרת לאפליקציה נכנסת לתור, מה שגורם לעלייה חדה בשיעור ההשהיה בכל המערכת. מכיוון שהערך המותאם אישית redis_client_rtt תואם לזמן האחזור של השרת (commands/usec_per_call), השרת שמריץ את הפקודה הוא צוואר הבקבוק.

פותחים את Trace ומסתכלים על פקודות Redis בטווחים האיטיים כדי לזהות איזו שאילתה גורמת לחסימה. מחליפים את פקודות החסימה בפקודות שלא חוסמות בקוד.

כדי לבצע איטרציה על מערכי נתונים גדולים באופן מצטבר בלי לנעול את שרשור השרת, משתמשים ב-SCAN, ב-SSCAN או ב-HSCAN.

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

  • workload.googleapis.com/redis_client_rtt (מדד בצד הלקוח): גבוה באופן עקבי (ערכי p50 ו-p99 גבוהים בערך ב-30 עד 100 אלפיות השנייה)
  • redis.googleapis.com/commands/
    usec_per_call
    (מדד של שרת Memorystore for Redis): נמוך מאוד (< 1ms)
  • workload.googleapis.com/
    redis_client_blocking_latency

    ו-redis_application_blocking_latency: נמוך / רגיל
שרת Redis מריץ פקודות באופן מיידי, אבל האפליקציה והמכונה שלכם נפרסים באזורים שונים (לדוגמה, us-central1 ו-us-east1). כל מנת נתונים ברשת צריכה לעבור בתשתית הענן הפיזית של Google Cloud בין מרכזי הנתונים הגיאוגרפיים האלה. התוצאה היא קנס חובה של מהירות האור בכל נסיעה הלוך ושוב בין אזורים. כדי להקטין את זמן האחזור, כדאי לפרוס את האפליקציה כך שהיא תמוקם באותו אזור ותחום כמו המופע. כדי לראות את האזור של האפליקציה והמופע, משתמשים במסוף Google Cloud .

המאמרים הבאים