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

如何用Cassandra与Go将秒级数据重采样为分钟级OHLC数据

问题:Cassandra分钟级OHLC重采样实现(Go语言)

源表结构

CREATE TABLE democassendra.tld2 (
id1 int,
id2 int,
year int,
month int,
day int,
hour int, 
min int, 
timestamp bigint,
price int,
PRIMARY KEY ((id1, id2, year, month, day, hour, min), timestamp)
)

该表每秒插入新数据,需每分钟重采样为分钟级数据。

目标表结构

CREATE TABLE democassendra.trd2 (
id1 int,
id2 int,
year int,
month int,
day int,
hour int, 
min int, 
open int,
high int,
low int,
close int,
PRIMARY KEY ((id1, id2, year, month, day, hour, min))
)

核心需求

每分钟执行以下操作:

  • 读取上一分钟的所有秒级数据
  • 按(id1, id2, year, month, day, hour, min)分组,计算每组的OHLC值:
    • open:组内第一条数据的price
    • high:组内price最大值
    • low:组内price最小值
    • close:组内最后一条数据的price
  • 将计算结果写入目标表trd2

示例数据

源表数据

id1    | id2    | year    | month    | day    | hour    | min    | timestamp (epoch)| price
-------------------------------------------------------------------------------------------
1      | 1      | 2023    | 7        | 20     | 15      | 2      | 123456           | 100
1      | 1      | 2023    | 7        | 20     | 15      | 2      | 123457           | 103
1      | 1      | 2023    | 7        | 20     | 15      | 2      | 123458           | 101
2      | 1      | 2023    | 7        | 20     | 15      | 2      | 123456           | 67
2      | 1      | 2023    | 7        | 20     | 15      | 2      | 123458           | 69

目标表数据

id1 | id2 | year | month | day | hour | min | OPEN| HIGH | LOW | CLOSE
-----------------------------------------------------------------------
1   | 1   | 2023 | 7     | 20  | 15   | 2   | 100 | 103  | 100 | 101 
2   | 1   | 2023 | 7     | 20  | 15   | 2   | 67  | 69   | 67  | 69

当前困境

已实现Cassandra连接、每分钟定时触发逻辑,但无法读取指定分钟的全量数据:Cassandra不支持仅用部分分区键(year, month, day, hour, min)过滤,必须提供全部分区键(id1, id2, year, month, day, hour, min)才能查询,导致无法批量获取指定分钟的所有分组数据。


解决方案

方案1:调整源表分区策略(推荐)

修改源表的分区键,将时间维度作为分区前缀,这样可以直接通过时间条件查询指定分钟的全量数据,调整后的表结构:

CREATE TABLE democassendra.tld2 (
id1 int,
id2 int,
year int,
month int,
day int,
hour int, 
min int, 
timestamp bigint,
price int,
PRIMARY KEY ((year, month, day, hour, min), id1, id2, timestamp)
)

调整后可直接按时间维度查询,再在Go代码中按id1, id2分组计算OHLC。

方案2:遍历小基数的id1, id2组合(仅适合小范围场景)

如果id1和id2的取值范围很小,可预先维护所有可能的组合,每分钟遍历这些组合,拼接全部分区键(id1, id2, 目标年/月/日/时/分)进行查询,逐一计算OHLC后写入目标表。

方案3:流处理工具辅助(大数据量场景)

若无法修改源表且id1, id2基数大,可使用Spark/Flink等流处理工具监听Cassandra CDC(变更日志),实时按分钟窗口聚合计算OHLC后写入目标表,适合高数据量场景。


Go代码实现示例(基于方案1的调整后表结构)

package main

import (
	"context"
	"fmt"
	"time"

	"github.com/gocql/gocql"
)

// 源表数据结构
type TLD2 struct {
	ID1       int
	ID2       int
	Timestamp int64
	Price     int
}

// 目标表数据结构
type TRD2 struct {
	ID1   int
	ID2   int
	Year  int
	Month int
	Day   int
	Hour  int
	Min   int
	Open  int
	High  int
	Low   int
	Close int
}

func main() {
	// 初始化Cassandra会话
	cluster := gocql.NewCluster("127.0.0.1")
	cluster.Keyspace = "democassendra"
	cluster.Consistency = gocql.Quorum
	session, err := cluster.CreateSession()
	if err != nil {
		panic(err)
	}
	defer session.Close()

	// 每分钟执行一次任务
	ticker := time.NewTicker(1 * time.Minute)
	defer ticker.Stop()

	for range ticker.C {
		// 获取上一分钟的时间参数
		prevMin := time.Now().Add(-1 * time.Minute).Truncate(time.Minute)
		year := prevMin.Year()
		month := int(prevMin.Month())
		day := prevMin.Day()
		hour := prevMin.Hour()
		min := prevMin.Minute()

		// 查询指定分钟的所有数据
		iter := session.Query(`SELECT id1, id2, timestamp, price FROM tld2 WHERE year = ? AND month = ? AND day = ? AND hour = ? AND min = ?`,
			year, month, day, hour, min).Iter()

		// 按(id1, id2)分组存储数据
		type groupKey struct {
			id1, id2 int
		}
		groups := make(map[groupKey][]TLD2)
		var row TLD2
		for iter.Scan(&row.ID1, &row.ID2, &row.Timestamp, &row.Price) {
			key := groupKey{row.ID1, row.ID2}
			groups[key] = append(groups[key], row)
		}
		if err := iter.Close(); err != nil {
			fmt.Printf("查询错误: %v\n", err)
			continue
		}

		// 计算OHLC并批量写入目标表
		batch := session.NewBatch(gocql.UnloggedBatch)
		for key, rows := range groups {
			if len(rows) == 0 {
				continue
			}
			// 利用Cassandra聚类键排序特性,直接取首尾值作为open/close
			open := rows[0].Price
			close := rows[len(rows)-1].Price
			high, low := open, open

			for _, r := range rows {
				if r.Price > high {
					high = r.Price
				}
				if r.Price < low {
					low = r.Price
				}
			}

			batch.Query(`INSERT INTO trd2 (id1, id2, year, month, day, hour, min, open, high, low, close) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
				key.id1, key.id2, year, month, day, hour, min, open, high, low, close)
		}

		if err := session.ExecuteBatch(context.Background(), batch); err != nil {
			fmt.Printf("批量写入错误: %v\n", err)
		} else {
			fmt.Printf("成功处理 %04d-%02d-%02d %02d:%02d 的数据\n", year, month, day, hour, min)
		}
	}
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 16:44:49