| title | Thread Safety |
|---|---|
| sidebar_position | 13 |
| id | thread-safety |
| license | Licensed to the Apache Software Foundation (ASF) under one or more contributor license agreements. See the NOTICE file distributed with this work for additional information regarding copyright ownership. The ASF licenses this file to You under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at http://www.apache.org/licenses/LICENSE-2.0 Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions and limitations under the License. |
This guide covers concurrent usage patterns for Fory Go, including the thread-safe wrapper and best practices for multi-goroutine environments.
The default Fory instance is not thread-safe:
f := fory.New(fory.WithXlang(true))
// NOT SAFE: Concurrent access from multiple goroutines
go func() {
f.Serialize(value1) // Race condition!
}()
go func() {
f.Serialize(value2) // Race condition!
}()For performance, Fory reuses internal state:
- Buffer is cleared and reused between calls
- Reference resolvers are reset
- Context objects are recycled
This avoids allocations but requires exclusive access.
For concurrent use, use the threadsafe package:
import "github.com/apache/fory/go/fory/threadsafe"
// Create thread-safe Fory
f := threadsafe.New()
// Safe for concurrent use
go func() {
data, _ := f.Serialize(value1)
}()
go func() {
data, _ := f.Serialize(value2)
}()The wrapper creates instances as needed and reuses them across goroutines. Each operation exclusively borrows one instance and returns it afterward. Registered types remain available when garbage collection reclaims cached instances. Serialized output is copied before returning, so callers can retain it safely.
// Create thread-safe instance
f := threadsafe.New()
// Instance methods
data, err := f.Serialize(value)
err = f.Deserialize(data, &target)
// Generic functions
data, err := threadsafe.Serialize(f, &value)
err = threadsafe.Deserialize(f, data, &target)
// Global convenience functions
data, err := threadsafe.Marshal(&value)
err = threadsafe.Unmarshal(data, &target)Register all types before the first serialization or deserialization operation. The first operation permanently freezes the wrapper's registrations, even if that operation fails. A later registration attempt returns an error.
f := threadsafe.New()
if err := f.RegisterStruct(User{}, 1); err != nil {
panic(err)
}
if err := f.RegisterStruct(Order{}, 2); err != nil {
panic(err)
}
// All concurrent operations use the registered types.
go func() {
data, err := f.Serialize(&User{ID: 1})
// Use data and handle err.
_, _ = data, err
}()RegisterStructByName, RegisterEnum, and RegisterEnumByName are also
available directly on the wrapper. Every registered type is available to all
concurrent operations and remains registered across garbage collections.
For custom per-instance initialization, use NewWithFactory:
f := threadsafe.NewWithFactory(func() *fory.Fory {
inner := fory.New()
if err := inner.RegisterExtension(CustomType{}, 100, newCustomSerializer()); err != nil {
panic(err)
}
return inner
})The factory may be called concurrently as additional instances are needed. It must return a fresh, identically configured instance on every call, with registrations completed before returning. Create stateful custom serializers separately for each instance.
With the default Fory, returned byte slices are views into the internal buffer:
f := fory.New(fory.WithXlang(true))
data1, _ := f.Serialize(value1)
// data1 is valid
data2, _ := f.Serialize(value2)
// data1 is NOW INVALID (buffer was reused)The thread-safe wrapper copies data automatically:
f := threadsafe.New()
data1, _ := f.Serialize(value1)
data2, _ := f.Serialize(value2)
// Both data1 and data2 are valid (independent copies)This is safer but has allocation overhead.
| Scenario | Non-Thread-Safe | Thread-Safe |
|---|---|---|
| Single goroutine | Fastest | Slower (pool overhead) |
| Multiple goroutines | Unsafe | Safe, good scaling |
| Memory allocations | Minimal | Per-call copy |
| Buffer reuse | Yes | Per-pool-instance |
func BenchmarkNonThreadSafe(b *testing.B) {
f := fory.New(fory.WithXlang(true))
f.RegisterStruct(User{}, 1)
user := &User{ID: 1, Name: "Alice"}
for i := 0; i < b.N; i++ {
data, _ := f.Serialize(user)
_ = data
}
}
func BenchmarkThreadSafe(b *testing.B) {
f := threadsafe.New()
f.RegisterStruct(User{}, 1)
user := &User{ID: 1, Name: "Alice"}
for i := 0; i < b.N; i++ {
data, _ := f.Serialize(user)
_ = data
}
}For maximum performance with known goroutine count:
func worker(id int) {
// Each worker has its own Fory instance
f := fory.New(fory.WithXlang(true))
f.RegisterStruct(User{}, 1)
for task := range tasks {
data, _ := f.Serialize(task)
process(data)
}
}
// Start workers
for i := 0; i < numWorkers; i++ {
go worker(i)
}For dynamic goroutine count or simplicity:
// Single shared instance
var f = threadsafe.New()
func init() {
f.RegisterStruct(User{}, 1)
}
func handleRequest(user *User) []byte {
// Safe from any goroutine
data, _ := f.Serialize(user)
return data
}var fory = threadsafe.New()
func init() {
fory.RegisterStruct(Response{}, 1)
}
func handler(w http.ResponseWriter, r *http.Request) {
response := &Response{
Status: "ok",
Data: getData(),
}
// Safe: threadsafe.Fory handles concurrency
data, err := fory.Serialize(response)
if err != nil {
http.Error(w, err.Error(), 500)
return
}
w.Header().Set("Content-Type", "application/octet-stream")
w.Write(data)
}// WRONG: Race condition
var f = fory.New(fory.WithXlang(true))
func handler1() {
f.Serialize(value1) // Race!
}
func handler2() {
f.Serialize(value2) // Race!
}Fix: Use threadsafe.New() or per-goroutine instances.
// WRONG: Buffer invalidated on next call
f := fory.New(fory.WithXlang(true))
data, _ := f.Serialize(value1)
savedData := data // Just copies the slice header!
f.Serialize(value2) // Invalidates data and savedDataFix: Clone the data or use thread-safe wrapper.
// Correct: Clone the data
data, _ := f.Serialize(value1)
savedData := make([]byte, len(data))
copy(savedData, data)
// Or use thread-safe (auto-copies)
f := threadsafe.New()
data, _ := f.Serialize(value1) // Already copied// RISKY: Concurrent registration
go func() {
f.RegisterStruct(TypeA{}, 1)
}()
go func() {
f.Serialize(value) // May not see TypeA
}()Fix: Register all types before the first serialization or deserialization.
- Register types at startup: Before the first serialization or deserialization
- Clone data if keeping references: With non-thread-safe instance
- Use per-worker instances for hot paths: Eliminates pool contention
- Profile before optimizing: Thread-safe overhead may be negligible