Skip to content
This repository was archived by the owner on Jul 28, 2026. It is now read-only.
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 19 additions & 6 deletions docs/configuration/tempo-config.md
Original file line number Diff line number Diff line change
Expand Up @@ -175,7 +175,7 @@ spanmetrics:
#
# In order to make a correct sampling decision it's important that the agent has
# a complete trace. This is achieved by waiting a given time for all the spans
# before evaluating the trace.
# before evaluating the trace. In order to get complete traces, configure group_by_trace.
#
# Tail sampling also supports multiple agent deployments, allowing to group all
# spans of a trace in the same agent by load balancing the spans by trace ID
Expand All @@ -186,11 +186,6 @@ tail_sampling:
policies:
- [<tailsamplingprocessor.policies>]

# Time that to wait before making a decision for a trace.
# Longer wait times reduce the probability of sampling an incomplete trace at
# the cost of higher memory usage.
decision_wait: [ <duration> | default="5s" ]

# load_balancing configures load balancing of spans across multiple agents.
# It ensures that all spans of a trace are sampled in the same instance.
# Only necessary if more than one agent process is receiving traces.
Expand Down Expand Up @@ -226,6 +221,24 @@ tail_sampling:
[ username: <string> ]
[ password: <secret> ]
[ password_file: <string> ]

# Configures aggregation of spans by trace.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We need better docs surrounding this. Perhaps we add some docs to the Tempo site and link them here? Things that need to be mentioned:

  1. Increased cpu/memory usage
  2. How rerouting the traces works (so teams can make sure that their network topologies are compatible)
  3. This only makes sense if you use tail sampling or service graphs.
  4. Explains the relationship between this and the load balancing processor.

# This is of particular interest for processor that benefit from having complete
# traces, such as tail-based sampling
#
# A trace will be considered complete after the defined `wait` duration has
# passed since its first span arrived to the processors.
# If the trace is not complete by then, it'll be split into more than one trace.
#
# Longer waiting times will increase the number of traces that are correctly
# grouped. However, it will also increase the memory overhead of the processor.
group_by_trace:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

shouldn't the "load balancing" settings be moved out as well? maybe renamed something to suggest it's routing by trace id.


# Defines the time to wait for a complete trace before considering it complete
[ wait: <duration> | default="5s" ]

# Configures the max amount of traces to keep in memory waiting for the duration
[ num_traces: <int> | default="1_000_000" ]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

max_traces?

```

> **Note:** More information on the following types can be found on the
Expand Down
30 changes: 30 additions & 0 deletions docs/upgrade-guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,36 @@ releases and how to migrate to newer versions.

These changes will come in a future version.

### Tempo: split grouping by trace from tail sampling config

Grouping spans by trace has been moved from an embedded functionality in tail
sampling to its own configuration block. This is done due to more processor
benefiting from receiving grouped traces, other than tail sampling.

As a consequence, `tail_sampling.wait_duration` has been deprecated in favor of
a `group_by_trace` block. Old configs will continue to work until it is fully
deprecated.

Example old config:

```yaml
tail_sampling:
duration_wait: 2s
policies:
- always_sample:
```

Example new config:

```yaml
tail_sampling:
policies:
- always_sample:
group_by_trace:
wait: 2s
```


### Logs: Deprecation of "loki" in config.

The term `loki` in the config has been deprecated of favor of `logs`. This
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -59,4 +59,6 @@ tempo:
dimensions:
- name: http.url
handler_endpoint: 0.0.0.0:8889
group_by_trace:
wait: 4s

4 changes: 2 additions & 2 deletions example/docker-compose/agent/config/agent.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ prometheus:
scrape_configs:
- job_name: local_scrape
static_configs:
- targets: ['127.0.0.1:12345']
- targets: ['127.0.0.1:12345', '127.0.0.1:8889']
labels:
cluster: 'docker_compose'
container: 'agent'
Expand Down Expand Up @@ -89,4 +89,4 @@ tempo:
processes: true
roots: true
spanmetrics:
prom_instance: tempo
handler_endpoint: 0.0.0.0:8889
1 change: 1 addition & 0 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ require (
github.com/olekukonko/tablewriter v0.0.2
github.com/oliver006/redis_exporter v1.15.0
github.com/open-telemetry/opentelemetry-collector-contrib/exporter/loadbalancingexporter v0.29.0
github.com/open-telemetry/opentelemetry-collector-contrib/processor/groupbytraceprocessor v0.29.0 // indirect
github.com/open-telemetry/opentelemetry-collector-contrib/processor/spanmetricsprocessor v0.29.0
github.com/open-telemetry/opentelemetry-collector-contrib/processor/tailsamplingprocessor v0.29.0
github.com/opentracing-contrib/go-grpc v0.0.0-20210225150812-73cb765af46e
Expand Down
2 changes: 2 additions & 0 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -1351,6 +1351,8 @@ github.com/open-telemetry/opentelemetry-collector-contrib/exporter/loadbalancing
github.com/open-telemetry/opentelemetry-collector-contrib/exporter/loadbalancingexporter v0.29.0/go.mod h1:6XSErVryxabB3T9aT7adMGoiPV7EG+hRACn2FRxQ/R4=
github.com/open-telemetry/opentelemetry-collector-contrib/pkg/batchpersignal v0.29.0 h1:zYbxW9OUwu/9W+P2bEOEqoQjn/rJfXmHxQQtYSpHLIc=
github.com/open-telemetry/opentelemetry-collector-contrib/pkg/batchpersignal v0.29.0/go.mod h1:zqMjJEuS55YssiJ4EFlD/dGBpZ/ij2jrtk0D2LC46wg=
github.com/open-telemetry/opentelemetry-collector-contrib/processor/groupbytraceprocessor v0.29.0 h1:ZgW0DOgUDFrRCPJ9pSUZ8iezpkUu/eQOP4RBDZINokI=
github.com/open-telemetry/opentelemetry-collector-contrib/processor/groupbytraceprocessor v0.29.0/go.mod h1:mEhEo1eK37dvXB3VzZEZxvQzjnuJw7jO6jnik7PcVtA=
github.com/open-telemetry/opentelemetry-collector-contrib/processor/spanmetricsprocessor v0.29.0 h1:0MwrtSxAQfk6DpktK8RK/yMADqmGZuK17F4U/zJnZTA=
github.com/open-telemetry/opentelemetry-collector-contrib/processor/spanmetricsprocessor v0.29.0/go.mod h1:6mho/Bo63I7dJqWyBvgfQSMObUopYisig9OBMDb33gw=
github.com/open-telemetry/opentelemetry-collector-contrib/processor/tailsamplingprocessor v0.29.0 h1:A10YxBMf8ki9bBNSbcG4bsnhhs7GXLegQ0xzjB5G40w=
Expand Down
5 changes: 3 additions & 2 deletions pkg/config/config_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -224,7 +224,7 @@ tempo:
backend: logs_instance
logs_instance_name: default
spans: true`,
expectedError: "error in config file: failed to validate automatic_logging for tempo config default: specified logs config default not found in agent config",
expectedError: "specified logs config default not found in agent config",
},
}

Expand All @@ -234,6 +234,7 @@ tempo:
return LoadBytes([]byte(tc.cfg), false, c)
})

require.EqualError(t, err, tc.expectedError)
require.Error(t, err)
require.Contains(t, err.Error(), tc.expectedError)
}
}
100 changes: 80 additions & 20 deletions pkg/tempo/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (
"fmt"
"io/ioutil"
"net"
"runtime"
"sort"
"time"

Expand All @@ -15,6 +16,7 @@ import (
"github.com/grafana/agent/pkg/tempo/promsdprocessor"
"github.com/grafana/agent/pkg/tempo/remotewriteexporter"
"github.com/open-telemetry/opentelemetry-collector-contrib/exporter/loadbalancingexporter"
"github.com/open-telemetry/opentelemetry-collector-contrib/processor/groupbytraceprocessor"
"github.com/open-telemetry/opentelemetry-collector-contrib/processor/spanmetricsprocessor"
"github.com/open-telemetry/opentelemetry-collector-contrib/processor/tailsamplingprocessor"
"github.com/prometheus/client_golang/prometheus"
Expand All @@ -38,8 +40,9 @@ import (
const (
spanMetricsPipelineName = "metrics/spanmetrics"

// defaultDecisionWait is the default time to wait for a trace before making a sampling decision
defaultDecisionWait = time.Second * 5
// defaultWaitDuration is the default time to wait for a trace before making a sampling decision
defaultWaitDuration = time.Second * 5
defaultNumTraces = 1_000_000

// defaultLoadBalancingPort is the default port the agent uses for internal load balancing
defaultLoadBalancingPort = "4318"
Expand Down Expand Up @@ -79,10 +82,8 @@ func (c *Config) Validate(logsConfig *logs.Config) error {
}

for _, inst := range c.Configs {
if inst.AutomaticLogging != nil {
if err := inst.AutomaticLogging.Validate(logsConfig); err != nil {
return fmt.Errorf("failed to validate automatic_logging for tempo config %s: %w", inst.Name, err)
}
if err := inst.Validate(logsConfig); err != nil {
return fmt.Errorf("failed validating config for tempo %s: %w", inst.Name, err)
}
}

Expand Down Expand Up @@ -118,7 +119,31 @@ type InstanceConfig struct {
AutomaticLogging *automaticloggingprocessor.AutomaticLoggingConfig `yaml:"automatic_logging,omitempty"`

// TailSampling defines a sampling strategy for the pipeline
TailSampling *tailSamplingConfig `yaml:"tail_sampling"`
TailSampling *tailSamplingConfig `yaml:"tail_sampling,omitempty"`

// GroupByTrace configures aggregation of spans by trace
// making processing of complete traces possible.
// This is useful for processing such as tail-based sampling.
GroupByTrace *groupByTraceConfig `yaml:"group_by_trace,omitempty"`
}

// Validate ensures that the InstanceConfig is valid
func (c *InstanceConfig) Validate(logsConfig *logs.Config) error {
if c.TailSampling != nil && c.GroupByTrace != nil {
if c.TailSampling.DecisionWait != 0 && c.GroupByTrace.WaitDuration != defaultWaitDuration {
return fmt.Errorf("must configure at most one of aggregate_by_trace.wait_duration and tail_sampling.decision_wait. tail_sampling.decision_wait is deprecated in favor of aggregate_by_trace.wait_duration")
}

c.GroupByTrace.WaitDuration, c.TailSampling.DecisionWait = c.TailSampling.DecisionWait, 0
}

if c.AutomaticLogging != nil {
if err := c.AutomaticLogging.Validate(logsConfig); err != nil {
return fmt.Errorf("failed to validate automatic_logging: %w", err)
}
}

return nil
}

const (
Expand Down Expand Up @@ -216,13 +241,21 @@ type tailSamplingConfig struct {
// For more information, refer to https://github.com/open-telemetry/opentelemetry-collector-contrib/tree/main/processor/tailsamplingprocessor
Policies []map[string]interface{} `yaml:"policies"`
// DecisionWait defines the time to wait for a complete trace before making a decision
// Deprecated
DecisionWait time.Duration `yaml:"decision_wait,omitempty"`
// Port is the port the instance will use to receive load balanced traces
Port string `yaml:"port"`
// LoadBalancing is used to distribute spans of the same trace to the same agent instance
LoadBalancing *loadBalancingConfig `yaml:"load_balancing"`
}

type groupByTraceConfig struct {
// WaitDuration defines the time to wait for a complete trace before considering it complete
WaitDuration time.Duration `yaml:"wait,omitempty"`
// NumTraces is the max number of traces to keep in memory waiting for the duration
NumTraces int `yaml:"num_traces,omitempty"`
}

// loadBalancingConfig defines the configuration for load balancing spans between agent instances
// loadBalancingConfig is an OTel exporter's config with extra resolver config
type loadBalancingConfig struct {
Expand Down Expand Up @@ -514,22 +547,15 @@ func (c *InstanceConfig) otelConfig() (*config.Config, error) {
}

if c.TailSampling != nil {
wait := defaultDecisionWait
if c.TailSampling.DecisionWait != 0 {
wait = c.TailSampling.DecisionWait
}

policies, err := formatPolicies(c.TailSampling.Policies)
if err != nil {
return nil, err
}

// tail_sampling should be executed before the batch processor
// TODO(mario.rodriguez): put attributes processor before tail_sampling. Maybe we want to sample on mutated spans
processorNames = append([]string{"tail_sampling"}, processorNames...)
processorNames = append(processorNames, "tail_sampling")
processors["tail_sampling"] = map[string]interface{}{
"policies": policies,
"decision_wait": wait,
"policies": policies,
}

if c.TailSampling.LoadBalancing != nil {
Expand All @@ -553,6 +579,37 @@ func (c *InstanceConfig) otelConfig() (*config.Config, error) {
}
}

groupByTrace := c.TailSampling != nil
if groupByTrace {
wait := defaultWaitDuration
numTraces := defaultNumTraces
if c.GroupByTrace != nil {
if c.GroupByTrace.WaitDuration > 0 {
wait = c.GroupByTrace.WaitDuration
}
if c.GroupByTrace.NumTraces > 0 {
numTraces = c.GroupByTrace.NumTraces
}
}

if c.TailSampling != nil {
tsp, ok := processors["tail_sampling"].(map[string]interface{})
if !ok {
return nil, fmt.Errorf("failed to configure tail sampling")
}
tsp["decision_wait"] = wait
tsp["num_traces"] = numTraces
processors["tail_sampling"] = tsp
} else {
processorNames = append(processorNames, "groupbytrace")
processors["groupbytrace"] = map[string]interface{}{
"wait_duration": wait,
"num_traces": numTraces,
"num_workers": runtime.NumCPU(),
}
}
}

// Build Pipelines
splitPipeline := c.TailSampling != nil && c.TailSampling.LoadBalancing != nil
orderedSplitProcessors := orderProcessors(processorNames, splitPipeline)
Expand Down Expand Up @@ -638,6 +695,7 @@ func tracingFactories() (component.Factories, error) {
}

processors, err := component.MakeProcessorFactoryMap(
groupbytraceprocessor.NewFactory(),
batchprocessor.NewFactory(),
attributesprocessor.NewFactory(),
promsdprocessor.NewFactory(),
Expand All @@ -664,9 +722,10 @@ func orderProcessors(processors []string, splitPipelines bool) [][]string {
order := map[string]int{
"attributes": 0,
"spanmetrics": 1,
"tail_sampling": 2,
"automatic_logging": 3,
"batch": 4,
"groupbytrace": 2,
"tail_sampling": 3,
"automatic_logging": 4,
"batch": 5,
}

sort.Slice(processors, func(i, j int) bool {
Expand All @@ -687,7 +746,8 @@ func orderProcessors(processors []string, splitPipelines bool) [][]string {
foundAt := len(processors)
for i, processor := range processors {
if processor == "batch" ||
processor == "tail_sampling" {
processor == "tail_sampling" ||
processor == "groupbytrace" {
foundAt = i
break
}
Expand Down
Loading