Modern Go Engineering
Chapitre 15
15 - Cloud avec Go
15 - Cloud avec Go
Cours 15 : Cloud avec Go
1. Docker Multi-stage pour Go
1.1 Dockerfile optimisé
# Stage 1: Build
FROM golang:1.22-alpine AS builder
RUN apk add --no-cache git ca-certificates
WORKDIR /app
COPY go.mod go.sum ./
RUN go mod download
COPY . .
RUN CGO_ENABLED=0 GOOS=linux GOARCH=amd64 \
go build -ldflags="-s -w" -o /app/service ./cmd/server
# Stage 2: Runtime (scratch)
FROM scratch
COPY --from=builder /etc/ssl/certs/ca-certificates.crt /etc/ssl/certs/
COPY --from=builder /app/service /service
EXPOSE 8080
ENTRYPOINT ["/service"]
# Alternative : distroless
FROM gcr.io/distroless/static-debian12:nonroot
COPY --from=builder /app/service /service
USER nonroot:nonroot
ENTRYPOINT ["/service"]
# Alternative : alpine (avec shell)
FROM alpine:3.19
RUN apk add --no-cache ca-certificates tzdata
COPY --from=builder /app/service /service
EXPOSE 8080
ENTRYPOINT ["/service"]
1.2 Docker compose
version: '3.8'
services:
app:
build:
context: .
dockerfile: Dockerfile
ports:
- "8080:8080"
environment:
- DB_HOST=postgres
- REDIS_HOST=redis
- KAFKA_BROKERS=kafka:9092
depends_on:
postgres:
condition: service_healthy
redis:
condition: service_started
healthcheck:
test: ["CMD", "/service", "health"]
interval: 30s
timeout: 10s
retries: 3
deploy:
resources:
limits:
cpus: '0.5'
memory: 256M
postgres:
image: postgres:16-alpine
environment:
POSTGRES_DB: myapp
POSTGRES_USER: app
POSTGRES_PASSWORD: secret
volumes:
- pgdata:/var/lib/postgresql/data
healthcheck:
test: ["CMD-SHELL", "pg_isready -U app"]
interval: 5s
redis:
image: redis:7-alpine
volumes:
- redisdata:/data
volumes:
pgdata:
redisdata:
2. AWS SDK v2
2.1 S3 (Stockage)
package s3
import (
"context"
"fmt"
"io"
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/service/s3"
)
type S3Client struct {
client *s3.Client
bucket string
}
func NewS3Client(ctx context.Context, bucket, region string) (*S3Client, error) {
cfg, err := config.LoadDefaultConfig(ctx, config.WithRegion(region))
if err != nil {
return nil, fmt.Errorf("load config: %w", err)
}
return &S3Client{
client: s3.NewFromConfig(cfg),
bucket: bucket,
}, nil
}
func (c *S3Client) Upload(ctx context.Context, key string, body io.Reader) error {
_, err := c.client.PutObject(ctx, &s3.PutObjectInput{
Bucket: aws.String(c.bucket),
Key: aws.String(key),
Body: body,
})
return err
}
func (c *S3Client) Download(ctx context.Context, key string) (io.ReadCloser, error) {
output, err := c.client.GetObject(ctx, &s3.GetObjectInput{
Bucket: aws.String(c.bucket),
Key: aws.String(key),
})
if err != nil {
return nil, err
}
return output.Body, nil
}
func (c *S3Client) Delete(ctx context.Context, key string) error {
_, err := c.client.DeleteObject(ctx, &s3.DeleteObjectInput{
Bucket: aws.String(c.bucket),
Key: aws.String(key),
})
return err
}
func (c *S3Client) List(ctx context.Context, prefix string) ([]string, error) {
var keys []string
paginator := s3.NewListObjectsV2Paginator(c.client, &s3.ListObjectsV2Input{
Bucket: aws.String(c.bucket),
Prefix: aws.String(prefix),
})
for paginator.HasMorePages() {
page, err := paginator.NextPage(ctx)
if err != nil {
return nil, err
}
for _, obj := range page.Contents {
keys = append(keys, *obj.Key)
}
}
return keys, nil
}
2.2 DynamoDB (NoSQL)
package dynamodb
import (
"context"
"fmt"
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/feature/dynamodb/attributevalue"
"github.com/aws/aws-sdk-go-v2/feature/dynamodb/expression"
"github.com/aws/aws-sdk-go-v2/service/dynamodb"
)
type User struct {
ID string `dynamodbav:"id"`
Name string `dynamodbav:"name"`
Email string `dynamodbav:"email"`
Roles []string `dynamodbav:"roles"`
}
type DynamoDBClient struct {
client *dynamodb.Client
table string
}
func New(ctx context.Context, table, region string) (*DynamoDBClient, error) {
cfg, err := config.LoadDefaultConfig(ctx, config.WithRegion(region))
if err != nil {
return nil, err
}
return &DynamoDBClient{
client: dynamodb.NewFromConfig(cfg),
table: table,
}, nil
}
func (d *DynamoDBClient) Put(ctx context.Context, user User) error {
item, err := attributevalue.MarshalMap(user)
if err != nil {
return err
}
_, err = d.client.PutItem(ctx, &dynamodb.PutItemInput{
TableName: aws.String(d.table),
Item: item,
})
return err
}
func (d *DynamoDBClient) Get(ctx context.Context, id string) (*User, error) {
result, err := d.client.GetItem(ctx, &dynamodb.GetItemInput{
TableName: aws.String(d.table),
Key: map[string]types.AttributeValue{
"id": &types.AttributeValueMemberS{Value: id},
},
})
if err != nil {
return nil, err
}
if result.Item == nil {
return nil, fmt.Errorf("not found")
}
var user User
if err := attributevalue.UnmarshalMap(result.Item, &user); err != nil {
return nil, err
}
return &user, nil
}
func (d *DynamoDBClient) Query(ctx context.Context, email string) ([]User, error) {
expr, err := expression.NewBuilder().
WithKeyCondition(expression.KeyEqual(expression.Key("email"), expression.Value(email))).
Build()
if err != nil {
return nil, err
}
result, err := d.client.Query(ctx, &dynamodb.QueryInput{
TableName: aws.String(d.table),
KeyConditionExpression: expr.KeyCondition(),
ExpressionAttributeNames: expr.Names(),
ExpressionAttributeValues: expr.Values(),
})
if err != nil {
return nil, err
}
var users []User
for _, item := range result.Items {
var user User
attributevalue.UnmarshalMap(item, &user)
users = append(users, user)
}
return users, nil
}
2.3 SQS (Queues)
package sqs
import (
"context"
"encoding/json"
"fmt"
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/aws/aws-sdk-go-v2/service/sqs/types"
)
type SQSClient struct {
client *sqs.Client
queueURL string
}
func New(ctx context.Context, queueURL, region string) (*SQSClient, error) {
cfg, err := config.LoadDefaultConfig(ctx, config.WithRegion(region))
if err != nil {
return nil, err
}
return &SQSClient{
client: sqs.NewFromConfig(cfg),
queueURL: queueURL,
}, nil
}
func (c *SQSClient) SendMessage(ctx context.Context, msg any) error {
body, err := json.Marshal(msg)
if err != nil {
return err
}
_, err = c.client.SendMessage(ctx, &sqs.SendMessageInput{
QueueUrl: aws.String(c.queueURL),
MessageBody: aws.String(string(body)),
})
return err
}
func (c *SQSClient) ReceiveMessages(ctx context.Context, maxMessages int32) ([]types.Message, error) {
result, err := c.client.ReceiveMessage(ctx, &sqs.ReceiveMessageInput{
QueueUrl: aws.String(c.queueURL),
MaxNumberOfMessages: maxMessages,
WaitTimeSeconds: 20, // Long polling
})
if err != nil {
return nil, err
}
return result.Messages, nil
}
func (c *SQSClient) DeleteMessage(ctx context.Context, receiptHandle string) error {
_, err := c.client.DeleteMessage(ctx, &sqs.DeleteMessageInput{
QueueUrl: aws.String(c.queueURL),
ReceiptHandle: aws.String(receiptHandle),
})
return err
}
2.4 Lambda
package lambda
import (
"context"
"encoding/json"
"log"
"github.com/aws/aws-lambda-go/events"
"github.com/aws/aws-lambda-go/lambda"
"github.com/aws/aws-sdk-go-v2/service/dynamodb"
)
type Event struct {
UserID string `json:"user_id"`
Action string `json:"action"`
}
type Handler struct {
db *dynamodb.Client
}
func NewHandler() *Handler {
return &Handler{}
}
func (h *Handler) HandleRequest(ctx context.Context, event json.RawMessage) (string, error) {
var e Event
if err := json.Unmarshal(event, &e); err != nil {
return "", err
}
log.Printf("Processing event: user=%s action=%s", e.UserID, e.Action)
return "OK", nil
}
// API Gateway event
func (h *Handler) HandleAPIGateway(ctx context.Context, event events.APIGatewayProxyRequest) (events.APIGatewayProxyResponse, error) {
return events.APIGatewayProxyResponse{
StatusCode: 200,
Headers: map[string]string{"Content-Type": "application/json"},
Body: `{"message": "Hello from Lambda"}`,
}, nil
}
// SQS event
func (h *Handler) HandleSQS(ctx context.Context, event events.SQSEvent) error {
for _, msg := range event.Records {
log.Printf("SQS message: %s", msg.Body)
}
return nil
}
func main() {
handler := NewHandler()
lambda.Start(handler.HandleRequest)
}
3. Kubernetes avec client-go
3.1 Informers et Watchers
package k8s
import (
"context"
"fmt"
"time"
appsv1 "k8s.io/api/apps/v1"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/informers"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/cache"
"k8s.io/client-go/tools/clientcmd"
)
type K8sClient struct {
clientset *kubernetes.Clientset
}
func NewInClusterClient() (*K8sClient, error) {
config, err := rest.InClusterConfig()
if err != nil {
return nil, fmt.Errorf("in-cluster config: %w", err)
}
clientset, err := kubernetes.NewForConfig(config)
if err != nil {
return nil, err
}
return &K8sClient{clientset: clientset}, nil
}
func NewOutOfClusterClient(kubeconfig string) (*K8sClient, error) {
config, err := clientcmd.BuildConfigFromFlags("", kubeconfig)
if err != nil {
return nil, err
}
clientset, err := kubernetes.NewForConfig(config)
if err != nil {
return nil, err
}
return &K8sClient{clientset: clientset}, nil
}
// Deployment operations
func (c *K8sClient) CreateDeployment(ctx context.Context, deploy *appsv1.Deployment) error {
_, err := c.clientset.AppsV1().Deployments(deploy.Namespace).Create(ctx, deploy, metav1.CreateOptions{})
return err
}
func (c *K8sClient) GetDeployment(ctx context.Context, namespace, name string) (*appsv1.Deployment, error) {
return c.clientset.AppsV1().Deployments(namespace).Get(ctx, name, metav1.GetOptions{})
}
func (c *K8sClient) UpdateDeployment(ctx context.Context, deploy *appsv1.Deployment) error {
_, err := c.clientset.AppsV1().Deployments(deploy.Namespace).Update(ctx, deploy, metav1.UpdateOptions{})
return err
}
func (c *K8sClient) DeleteDeployment(ctx context.Context, namespace, name string) error {
return c.clientset.AppsV1().Deployments(namespace).Delete(ctx, name, metav1.DeleteOptions{})
}
// Pod operations
func (c *K8sClient) ListPods(ctx context.Context, namespace string, labelSelector string) (*corev1.PodList, error) {
return c.clientset.CoreV1().Pods(namespace).List(ctx, metav1.ListOptions{
LabelSelector: labelSelector,
})
}
func (c *K8sClient) GetPodLogs(ctx context.Context, namespace, name string, tailLines int64) (string, error) {
req := c.clientset.CoreV1().Pods(namespace).GetLogs(name, &corev1.PodLogOptions{
TailLines: &tailLines,
})
data, err := req.DoRaw(ctx)
if err != nil {
return "", err
}
return string(data), nil
}
// Informer pattern (watch)
func (c *K8sClient) WatchPods(ctx context.Context, namespace string, handlers cache.ResourceEventHandler) {
factory := informers.NewSharedInformerFactoryWithOptions(
c.clientset,
10*time.Minute,
informers.WithNamespace(namespace),
)
informer := factory.Core().V1().Pods().Informer()
informer.AddEventHandler(handlers)
stop := make(chan struct{})
defer close(stop)
informer.Run(stop)
}
3.2 Controller-runtime (Operators)
package controllers
import (
"context"
"time"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/runtime"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/log"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
)
type PodReconciler struct {
client.Client
Scheme *runtime.Scheme
}
func (r *PodReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
log := log.FromContext(ctx)
var pod corev1.Pod
if err := r.Get(ctx, req.NamespacedName, &pod); err != nil {
return ctrl.Result{}, client.IgnoreNotFound(err)
}
log.Info("Reconciling Pod", "name", pod.Name, "phase", pod.Status.Phase)
// Handle pod phases
switch pod.Status.Phase {
case corev1.PodPending:
log.Info("Pod is pending")
case corev1.PodRunning:
log.Info("Pod is running")
case corev1.PodFailed:
log.Info("Pod failed", "reason", pod.Status.Reason)
}
return ctrl.Result{RequeueAfter: 30 * time.Second}, nil
}
func (r *PodReconciler) SetupWithManager(mgr ctrl.Manager) error {
return ctrl.NewControllerManagedBy(mgr).
For(&corev1.Pod{}).
Complete(r)
}
4. CloudEvents
package cloudevents
import (
"context"
"time"
ce "github.com/cloudevents/sdk-go/v2"
"github.com/cloudevents/sdk-go/v2/client"
"github.com/cloudevents/sdk-go/v2/protocol/http"
)
type EventPublisher struct {
client client.Client
}
func NewEventPublisher() (*EventPublisher, error) {
c, err := ce.NewClientHTTP()
if err != nil {
return nil, err
}
return &EventPublisher{client: c}, nil
}
func (p *EventPublisher) Publish(ctx context.Context, eventType, source string, data any) error {
event := ce.NewEvent()
event.SetID(GenerateID())
event.SetType(eventType)
event.SetSource(source)
event.SetTime(time.Now())
event.SetDataContentType("application/json")
if err := event.SetData("application/json", data); err != nil {
return err
}
ctx = ce.ContextWithTarget(ctx, "http://event-broker:8080")
return p.client.Send(ctx, event)
}
// Consumer
type EventConsumer struct {
client client.Client
}
func NewEventConsumer() (*EventConsumer, error) {
c, err := ce.NewClientHTTP(ce.WithPort(8080))
if err != nil {
return nil, err
}
return &EventConsumer{client: c}, nil
}
func (c *EventConsumer) Start(ctx context.Context, handler func(ce.Event)) {
c.client.StartReceiver(ctx, func(event ce.Event) {
handler(event)
})
}
5. Diagrammes
Diagramme en cours de génération...
Points Clés
- Docker multi-stage : scratch pour taille minimale, distroless pour sécurité
- AWS SDK v2 : modular, context-aware, paginators
- client-go : informers pour watch efficace, controller-runtime pour operators
- Lambda : handler simple pour API Gateway, SQS, S3 events
- CloudEvents : standard CNCF pour events cloud
- K8s operators : controller-runtime pattern