pool

package
v1.11.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jul 9, 2026 License: MIT Imports: 6 Imported by: 0

README

Connection Pool Package

Home  /  Connection Pool Package

 

The pool package provides connection pooling functionality for high-throughput RabbitMQ applications. Connection pooling allows applications to maintain multiple connections to RabbitMQ and distribute load across them, which is particularly useful for scenarios where a single connection might become a bottleneck.

 

Features

  • Round-robin connection selection - Distribute load evenly across connections
  • Automatic health monitoring - Periodic health checks with configurable intervals
  • Connection repair and recovery - Automatically recreate failed connections
  • Pool statistics and metrics - Monitor pool health and performance
  • Configurable behavior - Customize pool size, health check intervals, and repair settings
  • Thread-safe operations - Safe for concurrent use across goroutines

 

🔝 back to top

 

Quick Start

Basic Usage
package main

import (
    "context"
    "log"

    "github.com/cloudresty/go-rabbitmq"
    "github.com/cloudresty/go-rabbitmq/pool"
)

func main() {
    // Create a connection pool with 5 connections
    connectionPool, err := pool.New(5,
        pool.WithClientOptions(
            rabbitmq.WithHosts("localhost:5672"),
            rabbitmq.WithCredentials("guest", "guest"),
        ),
    )
    if err != nil {
        log.Fatal("Failed to create connection pool:", err)
    }
    defer connectionPool.Close()

    // Get a client from the pool
    client, err := connectionPool.Get()
    if err != nil {
        log.Fatal("No healthy connections available:", err)
    }

    // Use the client normally
    publisher, err := client.NewPublisher()
    if err != nil {
        log.Fatal("Failed to create publisher:", err)
    }
    defer publisher.Close()

    // Publish a message
    err = publisher.Publish(context.Background(), "test.queue", rabbitmq.Publishing{
        Body: []byte("Hello from pool!"),
    })
    if err != nil {
        log.Fatal("Failed to publish:", err)
    }
}

 

🔝 back to top

 

Advanced Configuration
// Create a pool with custom configuration using functional options
connectionPool, err := pool.New(10,
    pool.WithHealthCheck(30*time.Second),     // Health check frequency
    pool.WithAutoRepair(true),                // Auto-repair failed connections
    pool.WithClientOptions(
        rabbitmq.WithHosts("localhost:5672"),
        rabbitmq.WithCredentials("user", "pass"),
        rabbitmq.WithVHost("/prod"),
    ),
)
if err != nil {
    log.Fatal(err)
}
defer connectionPool.Close()

 

🔝 back to top

 

Pool Management

Getting Connections
// Get any available healthy client (fast, non-blocking)
client, err := connectionPool.Get()
if err != nil {
    log.Printf("No healthy connections available: %v", err)
    return
}

// Get a specific client by index (useful for sharding)
client := connectionPool.GetClientByIndex(2)

 

🔝 back to top

 

Pool Statistics
stats := connectionPool.GetStats()
fmt.Printf("Pool size: %d\n", stats.Size)
fmt.Printf("Healthy connections: %d\n", stats.HealthyConnections)
fmt.Printf("Unhealthy connections: %d\n", stats.UnhealthyConnections)
fmt.Printf("Total repair attempts: %d\n", stats.TotalRepairAttempts)
fmt.Printf("Last health check: %v\n", stats.LastHealthCheck)
fmt.Printf("Health monitoring enabled: %t\n", stats.HealthMonitoringEnabled)
fmt.Printf("Auto repair enabled: %t\n", stats.AutoRepairEnabled)

if len(stats.Errors) > 0 {
    fmt.Printf("Connection errors: %v\n", stats.Errors)
}

 

🔝 back to top

 

Configuration Options

The pool package uses functional options for configuration:

// WithHealthCheck enables health monitoring with specified interval
// Setting interval to 0 disables health monitoring
pool.WithHealthCheck(30 * time.Second)

// WithAutoRepair enables or disables automatic connection repair
pool.WithAutoRepair(true)

// WithClientOptions sets options for the RabbitMQ clients in the pool
pool.WithClientOptions(
    rabbitmq.WithHosts("localhost:5672"),
    rabbitmq.WithCredentials("user", "pass"),
    rabbitmq.WithConnectionName("my-pool"),
)

 

🔝 back to top

 

Health Monitoring

Health monitoring and connection repair are configured at pool creation time:

// Create pool with health monitoring disabled
pool, err := pool.New(5,
    pool.WithHealthCheck(0), // 0 disables health monitoring
    pool.WithClientOptions(rabbitmq.WithHosts("localhost:5672")),
)

// Create pool with custom health check interval
pool, err := pool.New(5,
    pool.WithHealthCheck(15*time.Second), // Check every 15 seconds
    pool.WithClientOptions(rabbitmq.WithHosts("localhost:5672")),
)

// Manual health check
ctx := context.Background()
err := connectionPool.HealthCheck(ctx)
if err != nil {
    log.Printf("Health check failed: %v", err)
}

 

🔝 back to top

 

Connection Repair

Automatic connection repair is also configured at creation time:

// Create pool with auto repair enabled (default)
pool, err := pool.New(5,
    pool.WithAutoRepair(true),
    pool.WithClientOptions(rabbitmq.WithHosts("localhost:5672")),
)

// Create pool with auto repair disabled
pool, err := pool.New(5,
    pool.WithAutoRepair(false),
    pool.WithClientOptions(rabbitmq.WithHosts("localhost:5672")),
)

 

🔝 back to top

 

Best Practices

Pool Sizing
  • Small pools (2-5 connections): For low-to-medium traffic applications
  • Medium pools (5-20 connections): For high-traffic applications
  • Large pools (20+ connections): For very high-throughput scenarios

 

// For most applications, start with 5 connections
pool, err := pool.New(5, opts...)

// Scale up based on your traffic patterns
pool, err := pool.New(20, opts...)  // High traffic

 

🔝 back to top

 

Health Monitoring Options
  • Development: Disable or use long intervals (5+ minutes)
  • Production: Use reasonable intervals (30-60 seconds)

 

// Development - minimal monitoring
config := pool.Config{
    Size:                5,
    HealthCheckInterval: 5 * time.Minute,  // Less frequent
    RepairEnabled:       false,            // Manual intervention
}

// Production - active monitoring
config := pool.Config{
    Size:                10,
    HealthCheckInterval: 30 * time.Second, // Regular checks
    RepairEnabled:       true,             // Auto-repair
    RepairThreshold:     10 * time.Second,
}

 

🔝 back to top

 

Error Handling
client, err := connectionPool.Get()
if err != nil {
    // No healthy connections available
    // Check pool stats for debugging
    stats := connectionPool.GetStats()
    log.Printf("Pool stats: %+v", stats)
    log.Printf("Error: %v", err)

    // Consider fallback strategy
    return handleNoConnectionsAvailable()
}

 

🔝 back to top

 

Resource Management
// Always close the pool when done
defer connectionPool.Close()

// In web servers, create pool once and reuse
var globalPool *pool.ConnectionPool

func init() {
    var err error
    globalPool, err = pool.New(10,
        rabbitmq.WithHosts("localhost:5672"),
    )
    if err != nil {
        log.Fatal(err)
    }
}

func handleRequest(w http.ResponseWriter, r *http.Request) {
    client, err := globalPool.Get()
    if err != nil {
        http.Error(w, "Service unavailable", http.StatusServiceUnavailable)
        return
    }
    // Use client...
}

 

🔝 back to top

 

Configuration Reference

Pool.Config
Field Type Description Default
Size int Number of connections in pool 5
MaxReconnectBackoff time.Duration Maximum backoff for reconnection 30s
HealthCheckInterval time.Duration Interval between health checks 30s
RepairEnabled bool Enable automatic connection repair true
RepairThreshold time.Duration Time before attempting repair 10s

 

🔝 back to top

 

Pool.Stats
Field Type Description
Size int Total pool size
HealthyConnections int Number of healthy connections
UnhealthyConnections int Number of unhealthy connections
Closed bool Whether pool is closed
Errors []string Current connection errors
TotalRepairAttempts int64 Total repair attempts made
LastHealthCheck time.Time Time of last health check
HealthMonitoringEnabled bool Health monitoring status
RepairEnabled bool Repair status

 

🔝 back to top

 

Integration Examples

With HTTP Server
func main() {
    // Create pool once
    connectionPool, err := pool.New(10,
        rabbitmq.WithHosts("localhost:5672"),
    )
    if err != nil {
        log.Fatal(err)
    }
    defer connectionPool.Close()

    // Use in handlers
    http.HandleFunc("/publish", func(w http.ResponseWriter, r *http.Request) {
        client, err := connectionPool.Get()
        if err != nil {
            http.Error(w, "No connections available", http.StatusServiceUnavailable)
            return
        }

        publisher, err := client.NewPublisher()
        if err != nil {
            http.Error(w, err.Error(), http.StatusInternalServerError)
            return
        }
        defer publisher.Close()

        // Publish message...
    })

    log.Fatal(http.ListenAndServe(":8080", nil))
}

 

🔝 back to top

 

With Worker Pool
func startWorkers(connectionPool *pool.ConnectionPool, numWorkers int) {
    for i := 0; i < numWorkers; i++ {
        go func(workerID int) {
            // Each worker gets its own client
            client := connectionPool.GetClientByIndex(workerID % connectionPool.Size())

            consumer, err := client.NewConsumer("task.queue")
            if err != nil {
                log.Printf("Worker %d failed to create consumer: %v", workerID, err)
                return
            }
            defer consumer.Close()

            // Consume messages...
        }(i)
    }
}

 

🔝 back to top

 

Performance Considerations

  • Connection overhead: Each connection uses memory and file descriptors
  • Health check cost: Frequent health checks add network overhead
  • Repair impact: Connection repairs cause temporary unavailability
  • Concurrent access: Pool operations are thread-safe but may block briefly

 

🔝 back to top

 

Testing

The pool package includes comprehensive tests. To run them:

go test ./pool

Note: Tests require a running RabbitMQ instance on localhost:5672.

 

🔝 back to top

 

The connection pool provides thread-safe access to multiple RabbitMQ connections with automatic health monitoring and repair capabilities.

 

🔝 back to top

 

 


Cloudresty

Website  |  LinkedIn  |  BlueSky  |  GitHub  |  Docker Hub

© Cloudresty - All rights reserved

 

Documentation

Overview

Package pool provides connection pooling functionality for high-throughput RabbitMQ applications.

Connection pooling allows applications to maintain multiple connections to RabbitMQ and distribute load across them. This is particularly useful for high-throughput scenarios where a single connection might become a bottleneck.

Features:

  • Round-robin connection selection
  • Automatic health monitoring
  • Connection repair and recovery
  • Pool statistics and metrics
  • Configurable pool size and behavior

Example usage:

// Create a connection pool with 5 connections
pool, err := pool.New(5, pool.WithClientOptions(rabbitmq.WithURL("amqp://localhost")))
if err != nil {
	log.Fatal(err)
}
defer pool.Close()

// Get a client from the pool
client, err := pool.Get()
if err != nil {
	log.Fatal("no healthy connections available:", err)
}

// Use the client normally
publisher, err := client.NewPublisher(...)

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type ConnectionPool

type ConnectionPool struct {
	// contains filtered or unexported fields
}

ConnectionPool manages multiple connections for high-throughput applications

func New

func New(size int, opts ...Option) (*ConnectionPool, error)

New creates a new connection pool with the specified size and options

func (*ConnectionPool) Close

func (p *ConnectionPool) Close() error

Close closes all connections in the pool

func (*ConnectionPool) Get

func (p *ConnectionPool) Get() (*rabbitmq.Client, error)

Get returns a healthy client from the pool (fast, non-blocking)

func (*ConnectionPool) GetClientByIndex

func (p *ConnectionPool) GetClientByIndex(index int) *rabbitmq.Client

GetClientByIndex returns a specific client by index (useful for sharding)

func (*ConnectionPool) GetStats

func (p *ConnectionPool) GetStats() Stats

GetStats returns statistics about the connection pool

func (*ConnectionPool) HealthCheck

func (p *ConnectionPool) HealthCheck(ctx context.Context) error

HealthCheck performs health checks on all connections in the pool

func (*ConnectionPool) HealthyCount

func (p *ConnectionPool) HealthyCount() int

HealthyCount returns the number of healthy connections

func (*ConnectionPool) Size

func (p *ConnectionPool) Size() int

Size returns the size of the connection pool

func (*ConnectionPool) Stats

func (p *ConnectionPool) Stats() rabbitmq.PoolStats

Stats converts internal stats to the interface PoolStats type

type Option

type Option func(*poolConfig)

Option configures the connection pool

func WithAutoRepair

func WithAutoRepair(enabled bool) Option

WithAutoRepair enables automatic connection repair

func WithClientOptions

func WithClientOptions(opts ...rabbitmq.Option) Option

WithClientOptions sets options for the RabbitMQ clients in the pool

func WithHealthCheck

func WithHealthCheck(interval time.Duration) Option

WithHealthCheck enables health monitoring with the specified interval

func WithLogger added in v1.1.3

func WithLogger(logger rabbitmq.Logger) Option

WithLogger sets the logger for the connection pool

type Stats

type Stats struct {
	Size                    int
	HealthyConnections      int
	UnhealthyConnections    int
	TotalRepairAttempts     int64
	LastHealthCheck         time.Time
	Closed                  bool
	Errors                  []string
	HealthMonitoringEnabled bool
	AutoRepairEnabled       bool
}

Stats contains statistics about the connection pool (public API)

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL