You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Go代码重构:将单一Schema的DataStore改造为支持多Schema写入

Go代码库重构:解耦DataStore与特定Schema,支持多数据集写入

问题背景

我正在维护一个Go代码库,公司有个叫DataStore的中央数据存储服务,能发布或读取不同数据集/Schema的数据。目前代码把特定数据集的Schema和DataStore实现硬耦合了,只能处理单一Schema(比如示例里的类似Order结构),现在需要重构代码,让它能支持写入任意Schema的数据集(比如Customers这类字段结构完全不同的)。

现有代码

type Parameters struct {
    SchemaName string
    Host       string
    Secret     string
}

type DataStore struct {
    params  *Parameters
    url     string
    request *http.Client
}

// 示例Schema(比如Order),其他Schema(如Customer)结构不同
type dataStoreSchema struct {
    eventUrl string `json:"url"`
    objectID string `json:"object_id"`
}

func New(parameters Parameters) (*DataStore, error) {
    // 初始化逻辑(比如拼接url、创建http.Client等)
    return &DataStore{
        params:  &parameters,
        url:     "https://" + parameters.Host + "/datastore",
        request: &http.Client{},
    }, nil
}

// 将输入转换为当前耦合的Schema对象
func (ds *DataStore) transformEvent(rawEvent interface{}) dataStoreSchema {
    // 假设这里是针对特定Schema的转换逻辑
    url := "example-url"
    objectId := "example-id"
    return dataStoreSchema{
        eventUrl: url,
        objectID: objectId,
    }
}

func (ds *DataStore) writeEvents(messages <-chan interface{}) error {
    for event := range messages {
        payload := ds.transformEvent(event)
        if err := ds.produce(payload); err != nil {
            return err
        }
    }
    return nil
}

func (ds *DataStore) produce(msg interface{}) error {
    // 实际写入DataStore的逻辑,比如序列化后发送HTTP请求
    jsonData, err := json.Marshal(msg)
    if err != nil {
        return err
    }
    _, err = ds.request.Post(ds.url, "application/json", bytes.NewBuffer(jsonData))
    return err
}

当前使用方式

params := Parameters{
    SchemaName: "Orders",
    Host:       "datastore.example.com",
    Secret:     "xxx",
}
myDatasetClient, _ := DataStore.New(params)
myDatasetClient.writeEvents(messageChan)

期望使用方式

希望能创建对应不同Schema的DataStore客户端,每个客户端绑定自己的转换逻辑:

// 针对Schema1的转换函数
func transformSchema1(raw interface{}) (interface{}, error) {
    // Schema1的转换逻辑
    return Schema1{/*...*/}, nil
}

// 针对Schema2的转换函数
func transformSchema2(raw interface{}) (interface{}, error) {
    // Schema2的转换逻辑
    return Schema2{/*...*/}, nil
}

params1 := Parameters{SchemaName: "Schema1", Host: "datastore.example.com", Secret: "xxx"}
myDatasetClient1, _ := DataStore.New(params1, transformSchema1)

params2 := Parameters{SchemaName: "Schema2", Host: "datastore.example.com", Secret: "xxx"}
myDatasetClient2, _ := DataStore.New(params2, transformSchema2)

重构方案

1. 抽象转换逻辑

定义一个函数类型,作为事件转换的抽象,让不同Schema可以传入自己的实现:

// EventTransformer 定义从原始事件到目标Schema对象的转换逻辑
type EventTransformer func(rawEvent interface{}) (interface{}, error)

2. 修改DataStore结构体

加入transformer字段,存储当前客户端对应的转换函数:

type DataStore struct {
    params      *Parameters
    url         string
    request     *http.Client
    transformer EventTransformer // 新增:绑定当前Schema的转换逻辑
}

3. 调整New函数

把转换函数作为参数传入,初始化DataStore时绑定:

func New(parameters Parameters, transformer EventTransformer) (*DataStore, error) {
    if transformer == nil {
        return nil, errors.New("transformer is required")
    }
    return &DataStore{
        params:      &parameters,
        url:         "https://" + parameters.Host + "/datastore",
        request:     &http.Client{},
        transformer: transformer,
    }, nil
}

4. 修改writeEvents方法

替换原来硬编码的transformEvent调用,使用绑定的transformer:

func (ds *DataStore) writeEvents(messages <-chan interface{}) error {
    for event := range messages {
        payload, err := ds.transformer(event)
        if err != nil {
            return fmt.Errorf("failed to transform event: %w", err)
        }
        if err := ds.produce(payload); err != nil {
            return fmt.Errorf("failed to produce event: %w", err)
        }
    }
    return nil
}

5. 适配produce方法

确保它能处理任意类型的payload(通过JSON序列化通用处理):

func (ds *DataStore) produce(msg interface{}) error {
    jsonData, err := json.Marshal(msg)
    if err != nil {
        return fmt.Errorf("failed to marshal payload: %w", err)
    }
    req, err := http.NewRequest("POST", ds.url, bytes.NewBuffer(jsonData))
    if err != nil {
        return err
    }
    req.Header.Set("Content-Type", "application/json")
    req.Header.Set("Authorization", "Bearer "+ds.params.Secret) // 假设需要Secret做鉴权
    _, err = ds.request.Do(req)
    return err
}

6. 定义各Schema结构体

比如针对不同业务的Schema:

// Schema1 比如Customers数据集的结构
type Schema1 struct {
    CustomerID string `json:"customer_id"`
    Email      string `json:"email"`
}

// Schema2 比如Orders数据集的结构
type Schema2 struct {
    OrderID    string `json:"order_id"`
    Amount     int    `json:"amount"`
    CustomerID string `json:"customer_id"`
}

重构后使用示例

// 实现Schema1的转换函数
func transformSchema1(raw interface{}) (interface{}, error) {
    // 假设原始事件是map类型,根据实际输入调整
    rawMap, ok := raw.(map[string]interface{})
    if !ok {
        return nil, errors.New("invalid raw event type for Schema1")
    }
    return Schema1{
        CustomerID: rawMap["id"].(string),
        Email:      rawMap["email"].(string),
    }, nil
}

// 实现Schema2的转换函数
func transformSchema2(raw interface{}) (interface{}, error) {
    rawMap, ok := raw.(map[string]interface{})
    if !ok {
        return nil, errors.New("invalid raw event type for Schema2")
    }
    return Schema2{
        OrderID:    rawMap["order_id"].(string),
        Amount:     int(rawMap["amount"].(float64)),
        CustomerID: rawMap["customer_id"].(string),
    }, nil
}

func main() {
    // 创建Schema1的客户端
    params1 := Parameters{
        SchemaName: "Customers",
        Host:       "datastore.example.com",
        Secret:     "my-secret-key",
    }
    client1, err := New(params1, transformSchema1)
    if err != nil {
        log.Fatal(err)
    }

    // 创建Schema2的客户端
    params2 := Parameters{
        SchemaName: "Orders",
        Host:       "datastore.example.com",
        Secret:     "my-secret-key",
    }
    client2, err := New(params2, transformSchema2)
    if err != nil {
        log.Fatal(err)
    }

    // 分别处理不同数据集的消息
    go client1.writeEvents(customerMessages)
    go client2.writeEvents(orderMessages)
}

这样就完全解耦了DataStore核心逻辑和具体的Schema转换,后续新增Schema只需要:

  • 定义对应的结构体
  • 实现转换函数
  • 创建新的DataStore客户端即可

内容的提问来源于stack exchange,提问作者wubbalubba

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.09 22:51:11