시간 기반 샤딩
시간 기반 샤딩 (tutorials-operations-time-based-sharding)
출처: Qdrant 공식문서
소셜 미디어나 이미지·비디오 스트림처럼 빠르게 유입되는 대용량 데이터셋을 다룰 때는 효율적인 저장과 검색이 매우 중요해요. 이런 데이터에서는 종종 최근 데이터만 관련이 있고, 오래된 데이터는 아카이브하거나 삭제할 수 있죠. 예를 들어 소셜 미디어 게시물의 감성 분석에서는 최신 트렌드를 파악하려고 최근 7일 치 데이터만 필요하고, 쿼리 대부분은 최근 24시간에 집중될 수 있어요.
모든 것을 기본 샤딩으로 Qdrant 컬렉션에 저장하면, 오래된 포인트를 삭제할 때 전체 데이터셋에 걸친 비싼 재인덱싱이 발생해서 성능에 영향을 줄 수 있어요. 더 나은 해결책은 시간 기반 샤딩(time-based sharding) 이에요. 타임스탬프를 기준으로 포인트를 특정 샤드(또는 샤드들)로 라우팅하는 방식이죠. 자연스러운 TTL(time-to-live) 구획이 있는 유스 케이스에서는 타임스탬프 기반 키로 샤딩함으로써 최근 데이터를 효율적으로 질의하고, 오래된 데이터는 매끄럽게 버릴 수 있어요.
예를 들어 일 단위 샤드로 구성하면 오늘 데이터는 오늘 샤드에, 어제 데이터는 어제 샤드에 저장돼요. 쿼리는 특정 샤드(예: 오늘 샤드)를 겨냥하거나, 날짜 범위를 덮도록 여러 샤드를 대상으로 할 수 있어요.

데이터 규모와 보존 요구 사항에 따라 시간별, 주별, 월별, 또는 유스 케이스에 맞는 다른 시간 간격으로 샤딩할 수 있어요.
이 튜토리얼은 시간 기반 샤딩을 구현하는 방법을 안내하며 다음을 다뤄요:
- 사용자 정의 샤딩을 가진 Qdrant 컬렉션 만들기
- 타임스탬프를 기준으로 과거 데이터를 올바른 샤드에 배치로 수집하기
- 새 데이터를 가장 최근 샤드에 할당하기
- 하나 이상의 샤드 질의하기
- 오래된 샤드 정리(pruning)하기
Qdrant 클라이언트 설치 및 초기화
먼저 Qdrant 클라이언트를 설치해요:
!pip install qdrant-client
다음으로 클라이언트를 초기화해요:
from qdrant_client import QdrantClient, models
client = QdrantClient(
url=QDRANT_URL,
api_key=QDRANT_API_KEY,
cloud_inference=True
)
const client = new QdrantClient({
url: QDRANT_URL,
apiKey: QDRANT_API_KEY,
});
let client = Qdrant::from_url(QDRANT_URL)
.api_key(QDRANT_API_KEY)
.build()?;
QdrantClient client =
new QdrantClient(
QdrantGrpcClient.newBuilder(QDRANT_URL, 6334, true)
.withApiKey(QDRANT_API_KEY)
.build());
var client = new QdrantClient(
host: QDRANT_URL,
https: true,
apiKey: QDRANT_API_KEY
);
client, err := qdrant.NewClient(&qdrant.Config{
Host: QDRANT_URL,
APIKey: QDRANT_API_KEY,
UseTLS: true,
})
이 튜토리얼은 Qdrant Cloud Inference를 사용해 벡터 임베딩을 생성한다고 가정해요. 임베딩 인프라를 직접 관리한다면 같은 원칙을 적용할 수 있지만, 코드 예시를 여러분의 임베딩 서비스에 맞게 바꿔야 해요.
컬렉션 만들기
사용자 정의 샤딩을 가진 컬렉션을 만들려면 샤딩 메서드를 custom으로 설정해요.
from qdrant_client import models
collection_name = "my_collection"
if client.collection_exists(collection_name=collection_name):
client.delete_collection(collection_name=collection_name)
client.create_collection(
collection_name=collection_name,
vectors_config={
"dense_vector": models.VectorParams(
size=384, distance=models.Distance.COSINE
)
},
sharding_method=models.ShardingMethod.CUSTOM
)
const collectionName = "my_collection";
if (await client.collectionExists(collectionName)) {
await client.deleteCollection(collectionName);
}
await client.createCollection(collectionName, {
vectors: {
dense_vector: {
size: 384,
distance: "Cosine",
},
},
sharding_method: "custom",
});
let collection_name = "my_collection";
if client.collection_exists(collection_name).await? {
client.delete_collection(collection_name).await?;
}
let mut vectors_config = VectorsConfigBuilder::default();
vectors_config.add_named_vector_params(
"dense_vector",
VectorParamsBuilder::new(384, Distance::Cosine),
);
client
.create_collection(
CreateCollectionBuilder::new(collection_name)
.vectors_config(vectors_config)
.sharding_method(ShardingMethod::Custom.into()),
)
.await?;
String collectionName = "my_collection";
if (client.collectionExistsAsync(collectionName).get()) {
client.deleteCollectionAsync(collectionName).get();
}
client.createCollectionAsync(
CreateCollection.newBuilder()
.setCollectionName(collectionName)
.setVectorsConfig(VectorsConfig.newBuilder().setParamsMap(
VectorParamsMap.newBuilder().putAllMap(Map.of(
"dense_vector",
VectorParams.newBuilder()
.setSize(384)
.setDistance(Distance.Cosine)
.build()))))
.setShardingMethod(ShardingMethod.Custom)
.build()
).get();
string collectionName = "my_collection";
if (await client.CollectionExistsAsync(collectionName))
await client.DeleteCollectionAsync(collectionName);
await client.CreateCollectionAsync(
collectionName: collectionName,
vectorsConfig: new VectorParamsMap
{
Map = {
["dense_vector"] = new VectorParams { Size = 384, Distance = Distance.Cosine }
}
},
shardingMethod: ShardingMethod.Custom
);
collectionName := "my_collection"
exists, err := client.CollectionExists(context.Background(), collectionName)
if exists {
client.DeleteCollection(context.Background(), collectionName)
}
client.CreateCollection(context.Background(), &qdrant.CreateCollection{
CollectionName: collectionName,
VectorsConfig: qdrant.NewVectorsConfigMap(
map[string]*qdrant.VectorParams{
"dense_vector": {
Size: 384,
Distance: qdrant.Distance_Cosine,
},
},
),
ShardingMethod: qdrant.ShardingMethod_Custom.Enum(),
})
커스텀 샤드는 샤드 키로 접근할 수 있어요. 이 튜토리얼에서 샤드 키는 각 데이터 포인트의 타임스탬프에서 추출한 YYYY-MM-DD 형식의 날짜예요.
이 컬렉션은 각 샤드 키마다 하나의 샤드를 가져요(데이터가 있는 각 날짜마다 별도의 샤드). 매우 큰 데이터셋의 경우 컬렉션에 shard_number를 구성하면 쓰기 처리량을 개선할 수 있어요. shard_number는 기본값이 1이에요. 더 높은 값으로 설정하면 샤드 키마다 여러 샤드를 만들어 클러스터의 여러 피어에 쓰기 부하를 분산할 수 있어요. 다만 샤드를 너무 많이 만들면 각 샤드가 리소스를 소비하고 오버헤드가 늘어나 성능 저하로 이어질 수 있으니 주의하세요. 데이터셋과 클러스터 구성에 최적인 샤드 수를 테스트해보세요.
사용자 정의 샤딩을 사용하는 컬렉션과 자동 샤딩을 사용하는 일반 컬렉션의 차이점 두 가지를 알아둘 필요가 있어요:
- 자동 샤딩을 쓰는 일반 컬렉션에서
shard_number는 컬렉션의 총 샤드 수를 결정해요. 그러나 사용자 정의 샤딩에서는shard_number가 샤드 키당 샤드 수를 결정해요. 컬렉션의 총 샤드 수는 샤드 키(날짜) 수에shard_number를 곱한 값과 같아요. - 일반 컬렉션에 적용할 수 있는 컬렉션 수준 구성 변경(예: HNSW 파라미터)은 사용자 정의 샤딩을 쓰는 컬렉션에도 적용할 수 있어요. 이런 변경은 기존 샤드와 앞으로 생성될 새 샤드 모두에 소급 적용돼요.
과거 데이터 수집
시계열 데이터는 종종 스트림으로 유입되고, 새 데이터 포인트가 계속 추가돼요. 각 데이터 포인트에는 생성 시각을 나타내는 타임스탬프가 포함될 수 있어요. 이 타임스탬프로 데이터 포인트가 어느 샤드에 속하는지 결정할 수 있어요. 데이터에 타임스탬프가 없다면 현재 시간을 사용할 수 있어요.
이 튜토리얼은 타임스탬프가 있는 소셜 미디어 게시물 샘플 데이터셋을 사용해요. 오늘이 2026년 4월 7일이라고 가정해볼게요. 먼저 이번 주(4월 1일~7일)의 과거 데이터 일부를 수집하는 것부터 시작해요. 데이터셋의 샘플 행 몇 개는 다음과 같아요:
| datetime | text |
|---|---|
| 2026-04-06T09:04:28 | April sunshine through the office window makes everything better. |
| 2026-04-06T09:04:32 | Morning stretch, good coffee, clear intentions. Monday: sorted. |
| 2026-04-06T09:05:52 | Grateful for a productive first day of the week. |
데이터셋을 업로드하고 각 날짜의 데이터를 자신의 샤드에 저장해요:
from qdrant_client.http.models import PointStruct, Document
import uuid
csv_url = 'https://raw.githubusercontent.com/qdrant/examples/refs/heads/master/time-based-sharding/social-media-posts.csv'
# Retrieve a list of existing shard keys in the collection
existing_shard_keys = list(client.list_shard_keys(collection_name=collection_name).shard_keys)
dense_model = "sentence-transformers/all-MiniLM-L6-v2"
batch_size = 100
current_date = None
buffer: list[PointStruct] = []
for row in parse_csv(csv_url):
shard_date = row['datetime'][:10] # Extract YYYY-MM-DD
if shard_date != current_date:
# Flush buffer for the previous date before switching
if buffer:
client.upload_points(
collection_name=collection_name,
points=buffer,
shard_key_selector=current_date,
)
buffer = []
# Create shard for the new date if it doesn't exist yet
if shard_date not in existing_shard_keys:
client.create_shard_key(collection_name, shard_date)
existing_shard_keys.append(shard_date)
current_date = shard_date
# Add point to buffer
buffer.append(PointStruct(
id=uuid.uuid4().hex,
payload={"text": row['text'], "datetime": row['datetime']},
vector={"dense_vector": Document(text=row["text"], model=dense_model)}
))
# Flush batch if buffer size exceeds batch size
if len(buffer) >= batch_size:
client.upload_points(
collection_name=collection_name,
points=buffer,
shard_key_selector=current_date,
)
buffer = []
# Flush remaining partial batch
if buffer:
client.upload_points(
collection_name=collection_name,
points=buffer,
shard_key_selector=current_date,
)
const csvUrl = "https://raw.githubusercontent.com/qdrant/examples/refs/heads/master/time-based-sharding/social-media-posts.csv";
// Retrieve a list of existing shard keys in the collection
const shardKeysResult = await client.listShardKeys(collectionName);
const existingShardKeys = new Set((shardKeysResult.shard_keys ?? []).map((d) => String(d.key)));
const denseModel = "sentence-transformers/all-MiniLM-L6-v2";
const batchSize = 100;
let currentDate = "";
let buffer: Extract<Parameters<typeof client.upsert>[1], { points: unknown }>['points'] = [];
for await (const { text: postText, datetime } of parseCSV(csvUrl)) {
const shardDate = datetime.slice(0, 10); // Extract YYYY-MM-DD
if (shardDate !== currentDate) {
// Flush buffer for the previous date before switching
if (buffer.length > 0) {
await client.upsert(collectionName, { points: buffer, shard_key: currentDate });
buffer = [];
}
// Create shard for the new date if it doesn't exist yet
if (!existingShardKeys.has(shardDate)) {
await client.createShardKey(collectionName, { shard_key: shardDate });
existingShardKeys.add(shardDate);
}
currentDate = shardDate;
}
// Add point to buffer
buffer.push({
id: crypto.randomUUID(),
vector: { dense_vector: { text: postText, model: denseModel } },
payload: { text: postText, datetime },
});
// Flush batch if buffer size exceeds batch size
if (buffer.length >= batchSize) {
await client.upsert(collectionName, { points: buffer, shard_key: currentDate });
buffer = [];
}
}
// Flush remaining partial batch
if (buffer.length > 0) {
await client.upsert(collectionName, { points: buffer, shard_key: currentDate });
}
let csv_url = "https://raw.githubusercontent.com/qdrant/examples/refs/heads/master/time-based-sharding/social-media-posts.csv";
// Retrieve a list of existing shard keys in the collection
let response = client.list_shard_keys(collection_name).await?;
let mut existing_shard_keys: HashSet<String> = response
.shard_keys
.into_iter()
.filter_map(|d| {
d.key?.key.and_then(|k| match k {
shard_key::Key::Keyword(s) => Some(s),
_ => None,
})
})
.collect();
let dense_model = "sentence-transformers/all-MiniLM-L6-v2";
let batch_size = 100;
let mut current_date = String::new();
let mut buffer: Vec<PointStruct> = Vec::new();
for row in parse_csv(csv_url)? {
let row = row?;
let text = row.text;
let datetime = row.datetime;
let shard_date = datetime[..10].to_string(); // Extract YYYY-MM-DD
if shard_date != current_date {
// Flush buffer for the previous date before switching
if !buffer.is_empty() {
client
.upsert_points(
UpsertPointsBuilder::new(
collection_name,
std::mem::take(&mut buffer),
)
.shard_key_selector(current_date.clone()),
)
.await?;
}
// Create shard for the new date if it doesn't exist yet
if !existing_shard_keys.contains(&shard_date) {
client
.create_shard_key(
CreateShardKeyRequestBuilder::new(collection_name).request(
CreateShardKeyBuilder::default().shard_key(shard_date.clone()),
),
)
.await?;
existing_shard_keys.insert(shard_date.clone());
}
current_date = shard_date;
}
// Add point to buffer
buffer.push(PointStruct::new(
uuid::Uuid::new_v4().to_string(),
HashMap::from([(
"dense_vector".to_string(),
DocumentBuilder::new(&text, dense_model).build(),
)]),
[("text", text.into()), ("datetime", datetime.into())],
));
// Flush batch if buffer size exceeds batch size
if buffer.len() >= batch_size {
client
.upsert_points(
UpsertPointsBuilder::new(collection_name, std::mem::take(&mut buffer))
.shard_key_selector(current_date.clone()),
)
.await?;
}
}
// Flush remaining partial batch
if !buffer.is_empty() {
client
.upsert_points(
UpsertPointsBuilder::new(collection_name, buffer)
.shard_key_selector(current_date.clone()),
)
.await?;
}
String csvUrl = "https://raw.githubusercontent.com/qdrant/examples/refs/heads/master/time-based-sharding/social-media-posts.csv";
// Retrieve a list of existing shard keys in the collection
var shardKeyDescriptions = client.listShardKeysAsync(collectionName).get();
Set<String> existingShardKeys = new HashSet<>();
for (var desc : shardKeyDescriptions) {
existingShardKeys.add(desc.getKey().getKeyword());
}
String denseModel = "sentence-transformers/all-MiniLM-L6-v2";
int batchSize = 100;
String currentDate = null;
List<PointStruct> buffer = new ArrayList<>();
try (var stream = parseCSV(csvUrl)) {
for (var row : (Iterable<CsvRow>) stream::iterator) {
String text = row.text;
String datetime = row.datetime;
String shardDate = datetime.substring(0, 10); // Extract YYYY-MM-DD
if (!shardDate.equals(currentDate)) {
// Flush buffer for the previous date before switching
if (!buffer.isEmpty()) {
client.upsertAsync(
UpsertPoints.newBuilder()
.setCollectionName(collectionName)
.addAllPoints(buffer)
.setShardKeySelector(shardKeySelector(currentDate))
.build()
).get();
buffer.clear();
}
// Create shard for the new date if it doesn't exist yet
if (!existingShardKeys.contains(shardDate)) {
client.createShardKeyAsync(
CreateShardKeyRequest.newBuilder()
.setCollectionName(collectionName)
.setRequest(CreateShardKey.newBuilder()
.setShardKey(shardKey(shardDate))
.build())
.build()
).get();
existingShardKeys.add(shardDate);
}
currentDate = shardDate;
}
// Add point to buffer
buffer.add(
PointStruct.newBuilder()
.setId(id(UUID.randomUUID()))
.setVectors(namedVectors(Map.of(
"dense_vector",
vector(Document.newBuilder()
.setText(text)
.setModel(denseModel)
.build()))))
.putAllPayload(Map.of("text", value(text), "datetime", value(datetime)))
.build());
// Flush batch if buffer size exceeds batch size
if (buffer.size() >= batchSize) {
client.upsertAsync(
UpsertPoints.newBuilder()
.setCollectionName(collectionName)
.addAllPoints(buffer)
.setShardKeySelector(shardKeySelector(currentDate))
.build()
).get();
buffer.clear();
}
}
}
// Flush remaining partial batch
if (!buffer.isEmpty()) {
client.upsertAsync(
UpsertPoints.newBuilder()
.setCollectionName(collectionName)
.addAllPoints(buffer)
.setShardKeySelector(shardKeySelector(currentDate))
.build()
).get();
}
string csvUrl = "https://raw.githubusercontent.com/qdrant/examples/refs/heads/master/time-based-sharding/social-media-posts.csv";
// Retrieve a list of existing shard keys in the collection
var existingShardKeys = (await client.ListShardKeysAsync(collectionName))
.Select(sk => sk.Key.Keyword)
.ToHashSet();
string denseModel = "sentence-transformers/all-MiniLM-L6-v2";
int batchSize = 100;
string? currentDate = null;
var buffer = new List<PointStruct>();
await foreach (var (text, datetime) in ParseCsv(csvUrl))
{
string shardDate = datetime[..10]; // Extract YYYY-MM-DD
if (shardDate != currentDate)
{
// Flush buffer for the previous date before switching
if (buffer.Count > 0)
{
await client.UpsertAsync(
collectionName: collectionName,
points: buffer,
shardKeySelector: new ShardKeySelector
{
ShardKeys = { new List<ShardKey> { currentDate! } }
}
);
buffer.Clear();
}
// Create shard for the new date if it doesn't exist yet
if (!existingShardKeys.Contains(shardDate))
{
await client.CreateShardKeyAsync(
collectionName,
new CreateShardKey { ShardKey = new ShardKey { Keyword = shardDate } }
);
existingShardKeys.Add(shardDate);
}
currentDate = shardDate;
}
// Add point to buffer
buffer.Add(new PointStruct
{
Id = Guid.NewGuid(),
Vectors = new Dictionary<string, Vector>
{
["dense_vector"] = new Document { Text = text, Model = denseModel }
},
Payload = { ["text"] = text, ["datetime"] = datetime }
});
// Flush batch if buffer size exceeds batch size
if (buffer.Count >= batchSize)
{
await client.UpsertAsync(
collectionName: collectionName,
points: buffer,
shardKeySelector: new ShardKeySelector
{
ShardKeys = { new List<ShardKey> { currentDate! } }
}
);
buffer.Clear();
}
}
// Flush remaining partial batch
if (buffer.Count > 0)
{
await client.UpsertAsync(
collectionName: collectionName,
points: buffer,
shardKeySelector: new ShardKeySelector
{
ShardKeys = { new List<ShardKey> { currentDate! } }
}
);
}
csvUrl := "https://raw.githubusercontent.com/qdrant/examples/refs/heads/master/time-based-sharding/social-media-posts.csv"
shardKeyDescriptions, err := client.ListShardKeys(context.Background(), collectionName)
// Retrieve a list of existing shard keys in the collection
existingShardKeys := make(map[string]bool)
for _, desc := range shardKeyDescriptions {
existingShardKeys[desc.Key.GetKeyword()] = true
}
denseModel := "sentence-transformers/all-MiniLM-L6-v2"
batchSize := 100
var currentDate string
var buffer []*qdrant.PointStruct
err = parseCSV(csvUrl, func(row CSVRow) {
text := row.Text
datetime := row.Datetime
shardDate := datetime[:10] // Extract YYYY-MM-DD
if shardDate != currentDate {
// Flush buffer for the previous date before switching
if len(buffer) > 0 {
client.Upsert(context.Background(), &qdrant.UpsertPoints{
CollectionName: collectionName,
Points: buffer,
ShardKeySelector: &qdrant.ShardKeySelector{
ShardKeys: []*qdrant.ShardKey{qdrant.NewShardKey(currentDate)},
},
})
buffer = nil
}
// Create shard for the new date if it doesn't exist yet
if !existingShardKeys[shardDate] {
client.CreateShardKey(context.Background(), collectionName, &qdrant.CreateShardKey{
ShardKey: qdrant.NewShardKey(shardDate),
})
existingShardKeys[shardDate] = true
}
currentDate = shardDate
}
// Add point to buffer
buffer = append(buffer, &qdrant.PointStruct{
Id: qdrant.NewID(uuid.New().String()),
Vectors: qdrant.NewVectorsMap(map[string]*qdrant.Vector{
"dense_vector": qdrant.NewVectorDocument(&qdrant.Document{
Text: text,
Model: denseModel,
}),
}),
Payload: qdrant.NewValueMap(map[string]any{
"text": text,
"datetime": datetime,
}),
})
// Flush batch if buffer size exceeds batch size
if len(buffer) >= batchSize {
client.Upsert(context.Background(), &qdrant.UpsertPoints{
CollectionName: collectionName,
Points: buffer,
ShardKeySelector: &qdrant.ShardKeySelector{
ShardKeys: []*qdrant.ShardKey{qdrant.NewShardKey(currentDate)},
},
})
buffer = nil
}
})
// Flush remaining partial batch
if len(buffer) > 0 {
client.Upsert(context.Background(), &qdrant.UpsertPoints{
CollectionName: collectionName,
Points: buffer,
ShardKeySelector: &qdrant.ShardKeySelector{
ShardKeys: []*qdrant.ShardKey{qdrant.NewShardKey(currentDate)},
},
})
}
코드를 하나씩 살펴볼게요:
- 먼저 컬렉션에 존재하는 샤드 키 목록을 조회해요. 방금 컬렉션을 만들었으니 없어야 하지만, 운영 환경에서는 기존 데이터를 고려하도록 보장해주죠.
- 다음으로 헬퍼 함수로 URL에서 CSV 파일을 스트리밍해 파싱해요. 상세:
import csv
import urllib.request
def parse_csv(url):
with urllib.request.urlopen(url) as response:
reader = csv.DictReader(line.decode('utf-8') for line in response)
yield from reader
function parseCsvLine(line: string): string[] {
const fields: string[] = [];
let i = 0;
while (i < line.length) {
if (line[i] === '"') {
i++;
let field = "";
while (i < line.length) {
if (line[i] === '"' && line[i + 1] === '"') { field += '"'; i += 2; }
else if (line[i] === '"') { i++; break; }
else { field += line[i++]; }
}
fields.push(field);
if (line[i] === ",") i++;
} else {
const start = i;
while (i < line.length && line[i] !== ",") i++;
fields.push(line.slice(start, i));
if (i < line.length) i++;
}
}
return fields;
}
async function* parseCSV(url: string): AsyncGenerator<{ text: string; datetime: string }> {
const response = await fetch(url);
const reader = response.body!.getReader();
const decoder = new TextDecoder();
let remainder = "";
let headers: string[] | null = null;
let textIdx = -1;
let datetimeIdx = -1;
while (true) {
const { done, value } = await reader.read();
const chunk = done ? "" : decoder.decode(value, { stream: true });
const lines = (remainder + chunk).split("\n");
remainder = done ? "" : lines.pop()!;
for (const line of lines) {
if (!line.trim()) continue;
if (headers === null) {
headers = line.split(",");
textIdx = headers.indexOf("text");
datetimeIdx = headers.indexOf("datetime");
continue;
}
const fields = parseCsvLine(line);
yield { text: fields[textIdx], datetime: fields[datetimeIdx] };
}
if (done) break;
}
}
struct CsvRow {
text: String,
datetime: String,
}
fn parse_csv(url: &str) -> anyhow::Result<impl Iterator<Item = anyhow::Result<CsvRow>>> {
let reader = ureq::get(url).call()?.into_body().into_reader();
let mut rdr = csv::Reader::from_reader(reader);
let headers = rdr.headers()?.clone();
let text_idx = headers.iter().position(|h| h == "text").unwrap();
let datetime_idx = headers.iter().position(|h| h == "datetime").unwrap();
let iter = rdr.into_records().map(move |result| {
let record = result?;
Ok(CsvRow {
text: record[text_idx].to_string(),
datetime: record[datetime_idx].to_string(),
})
});
Ok(iter)
}
static class CsvRow {
final String text;
final String datetime;
CsvRow(String text, String datetime) { this.text = text; this.datetime = datetime; }
}
static Stream<CsvRow> parseCSV(String url) throws Exception {
Function<String, List<String>> parseCsvLine = line -> {
List<String> fields = new ArrayList<>();
boolean inQuotes = false;
var sb = new StringBuilder();
for (char c : line.toCharArray()) {
if (c == '"') {
inQuotes = !inQuotes;
} else if (c == ',' && !inQuotes) {
fields.add(sb.toString());
sb.setLength(0);
} else {
sb.append(c);
}
}
fields.add(sb.toString());
return fields;
};
var reader = new BufferedReader(new InputStreamReader(new URL(url).openStream()));
String headerLine = reader.readLine();
List<String> headers = List.of(headerLine.split(","));
int textIdx = headers.indexOf("text");
int datetimeIdx = headers.indexOf("datetime");
return reader.lines()
.map(line -> {
List<String> fields = parseCsvLine.apply(line);
return new CsvRow(fields.get(textIdx), fields.get(datetimeIdx));
})
.onClose(() -> { try { reader.close(); } catch (Exception ignored) {} });
}
async IAsyncEnumerable<(string text, string datetime)> ParseCsv(string url)
{
using var httpClient = new HttpClient();
using var stream = await httpClient.GetStreamAsync(url);
using var parser = new TextFieldParser(new StreamReader(stream));
parser.TextFieldType = Microsoft.VisualBasic.FileIO.FieldType.Delimited;
parser.SetDelimiters(",");
string[]? headers = parser.ReadFields();
int textIdx = Array.IndexOf(headers!, "text");
int datetimeIdx = Array.IndexOf(headers!, "datetime");
while (!parser.EndOfData)
{
var fields = parser.ReadFields()!;
yield return (fields[textIdx], fields[datetimeIdx]);
}
}
type CSVRow struct {
Text string
Datetime string
}
func parseCSV(url string, fn func(CSVRow)) error {
resp, err := http.Get(url)
if err != nil {
return err
}
defer resp.Body.Close()
csvReader := csv.NewReader(resp.Body)
headers, err := csvReader.Read()
if err != nil {
return err
}
textIdx, datetimeIdx := -1, -1
for i, h := range headers {
switch h {
case "text":
textIdx = i
case "datetime":
datetimeIdx = i
}
}
for {
row, err := csvReader.Read()
if err == io.EOF {
break
}
if err != nil {
return err
}
fn(CSVRow{Text: row[textIdx], Datetime: row[datetimeIdx]})
}
return nil
}
- CSV 파일은 행 단위로 스트리밍되며, 효율적인 업로드를 위해 포인트를 100개 단위 배치로 버퍼링해요. 최적의 배치 크기는 데이터와 클러스터에 따라 달라지므로, 최고 성능을 위해 여러 크기를 실험해보는 것이 좋아요.
- 각 행의 datetime 필드에서 날짜(
YYYY-MM-DD)를 추출해요. 이 날짜가 샤드 키로 사용돼 데이터를 올바른 샤드로 라우팅해요. - 마주치는 각 새 날짜에 대해, 아직 없다면 새 샤드를 만들어요.
- 스트림 중간에 날짜가 바뀌면 버퍼가 이전 날짜의 샤드로 플러시돼, 게시물이 잘못된 샤드에 쓰이지 않도록 해요.
- 데이터가 쓰이며 각 포인트는 다음을 받아요:
- 포인트 ID로 임의의 UUID
- 페이로드로 게시물
text와datetime sentence-transformers/all-MiniLM-L6-v2로 게시물 텍스트에서 생성한 dense 벡터 임베딩
- 100개 포인트의 전체 배치가 각각 shard key selector 파라미터로 올바른 날짜 기반 샤드를 대상으로 Qdrant에 업로드돼요.
- 루프가 끝난 뒤 남은 포인트는 부분적인 마지막 배치로 업로드돼요.
데이터 질의하기
오늘 데이터 질의
이제 게시물에 대한 의미론적 쿼리를 실행할 수 있어요. shard key selector를 2026-04-07로 설정하면(오늘이 2026년 4월 7일이라고 가정) 쿼리가 오늘 데이터를 저장하는 샤드로 제한돼요:
query_text = "coffee"
resp = client.query_points(
collection_name=collection_name,
query=Document(text=query_text, model=dense_model),
using="dense_vector",
limit=5,
shard_key_selector="2026-04-07"
)
print(resp)
const queryText = "coffee";
const singleShardResult = await client.query(collectionName, {
query: { text: queryText, model: denseModel },
using: "dense_vector",
limit: 5,
shard_key: "2026-04-07",
});
for (const hit of singleShardResult.points) {
console.log(hit);
}
let query_text = "coffee";
let result = client
.query(
QueryPointsBuilder::new(collection_name)
.query(Query::new_nearest(Document::new(query_text, dense_model)))
.using("dense_vector")
.limit(5)
.shard_key_selector("2026-04-07".to_string()),
)
.await?;
for hit in result.result {
println!("{:?}", hit);
}
String queryText = "coffee";
var result = client.queryAsync(
QueryPoints.newBuilder()
.setCollectionName(collectionName)
.setQuery(nearest(Document.newBuilder().setText(queryText).setModel(denseModel).build()))
.setUsing("dense_vector")
.setLimit(5)
.setShardKeySelector(shardKeySelector("2026-04-07"))
.build()
).get();
for (var hit : result) {
System.out.println(hit);
}
string queryText = "coffee";
var result = await client.QueryAsync(
collectionName: collectionName,
query: new Document { Text = queryText, Model = denseModel },
usingVector: "dense_vector",
limit: 5,
shardKeySelector: new ShardKeySelector
{
ShardKeys = { new List<ShardKey> { "2026-04-07" } }
}
);
foreach (var hit in result)
Console.WriteLine(hit);
queryText := "coffee"
result, err := client.Query(context.Background(), &qdrant.QueryPoints{
CollectionName: collectionName,
Query: qdrant.NewQueryDocument(&qdrant.Document{Text: queryText, Model: denseModel}),
Using: qdrant.PtrOf("dense_vector"),
Limit: qdrant.PtrOf(uint64(5)),
ShardKeySelector: &qdrant.ShardKeySelector{
ShardKeys: []*qdrant.ShardKey{qdrant.NewShardKey("2026-04-07")},
},
})
for _, hit := range result {
fmt.Println(hit)
}
여러 날짜의 데이터 질의
여러 샤드를 질의하려면 shard key selector를 샤드 키 목록으로 설정해요. 예를 들어 최근 2일(4월 6~7일)의 데이터를 질의하려면:
resp = client.query_points(
collection_name=collection_name,
query=Document(text=query_text, model=dense_model),
using="dense_vector",
limit=5,
shard_key_selector=["2026-04-06","2026-04-07"]
)
print(resp)
const multiShardResult = await client.query(collectionName, {
query: { text: queryText, model: denseModel },
using: "dense_vector",
limit: 5,
shard_key: ["2026-04-06", "2026-04-07"],
});
for (const hit of multiShardResult.points) {
console.log(hit);
}
let result = client
.query(
QueryPointsBuilder::new(collection_name)
.query(Query::new_nearest(Document::new(query_text, dense_model)))
.using("dense_vector")
.limit(5)
.shard_key_selector(ShardKeySelector {
shard_keys: vec![
"2026-04-06".to_string().into(),
"2026-04-07".to_string().into(),
],
fallback: None,
}),
)
.await?;
for hit in result.result {
println!("{:?}", hit);
}
result = client.queryAsync(
QueryPoints.newBuilder()
.setCollectionName(collectionName)
.setQuery(nearest(Document.newBuilder().setText(queryText).setModel(denseModel).build()))
.setUsing("dense_vector")
.setLimit(5)
.setShardKeySelector(ShardKeySelector.newBuilder()
.addShardKeys(shardKey("2026-04-06"))
.addShardKeys(shardKey("2026-04-07"))
.build())
.build()
).get();
for (var hit : result) {
System.out.println(hit);
}
result = await client.QueryAsync(
collectionName: collectionName,
query: new Document { Text = queryText, Model = denseModel },
usingVector: "dense_vector",
limit: 5,
shardKeySelector: new ShardKeySelector
{
ShardKeys = { new List<ShardKey> { "2026-04-06", "2026-04-07" } }
}
);
foreach (var hit in result)
Console.WriteLine(hit);
result, err = client.Query(context.Background(), &qdrant.QueryPoints{
CollectionName: collectionName,
Query: qdrant.NewQueryDocument(&qdrant.Document{Text: queryText, Model: denseModel}),
Using: qdrant.PtrOf("dense_vector"),
Limit: qdrant.PtrOf(uint64(5)),
ShardKeySelector: &qdrant.ShardKeySelector{
ShardKeys: []*qdrant.ShardKey{
qdrant.NewShardKey("2026-04-06"),
qdrant.NewShardKey("2026-04-07"),
},
},
})
for _, hit := range result {
fmt.Println(hit)
}
전체 데이터셋 질의
전체 데이터셋(모든 샤드)을 질의하려면 shard key selector 파라미터를 생략해요:
resp = client.query_points(
collection_name=collection_name,
query=Document(text=query_text, model=dense_model),
using="dense_vector",
limit=5,
)
print(resp)
const allShardsResult = await client.query(collectionName, {
query: { text: queryText, model: denseModel },
using: "dense_vector",
limit: 5,
});
for (const hit of allShardsResult.points) {
console.log(hit);
}
let result = client
.query(
QueryPointsBuilder::new(collection_name)
.query(Query::new_nearest(Document::new(query_text, dense_model)))
.using("dense_vector")
.limit(5),
)
.await?;
for hit in result.result {
println!("{:?}", hit);
}
result = client.queryAsync(
QueryPoints.newBuilder()
.setCollectionName(collectionName)
.setQuery(nearest(Document.newBuilder().setText(queryText).setModel(denseModel).build()))
.setUsing("dense_vector")
.setLimit(5)
.build()
).get();
for (var hit : result) {
System.out.println(hit);
}
result = await client.QueryAsync(
collectionName: collectionName,
query: new Document { Text = queryText, Model = denseModel },
usingVector: "dense_vector",
limit: 5
);
foreach (var hit in result)
Console.WriteLine(hit);
result, err = client.Query(context.Background(), &qdrant.QueryPoints{
CollectionName: collectionName,
Query: qdrant.NewQueryDocument(&qdrant.Document{Text: queryText, Model: denseModel}),
Using: qdrant.PtrOf("dense_vector"),
Limit: qdrant.PtrOf(uint64(5)),
})
for _, hit := range result {
fmt.Println(hit)
}
매일 자정에, 그날 수집될 새 데이터를 위한 새 샤드를 만들어요. 최근 7일 치 데이터만 질의한다면 가장 오래된 샤드를 삭제할 수도 있어요. 이것은 cron 작업으로 자동화할 수 있어요.
from datetime import date, timedelta
today = "2026-04-08"
oldest_shard_key = (date.fromisoformat(today) - timedelta(days=7)).isoformat()
client.create_shard_key(collection_name, today)
client.delete_shard_key(collection_name, oldest_shard_key)
const today = "2026-04-08";
const oldestDate = new Date(today);
oldestDate.setDate(oldestDate.getDate() - 7);
const oldestShardKey = oldestDate.toISOString().slice(0, 10);
await client.createShardKey(collectionName, { shard_key: today });
await client.deleteShardKey(collectionName, { shard_key: oldestShardKey });
let today = "2026-04-08";
let oldest_shard_key = (NaiveDate::parse_from_str(today, "%Y-%m-%d")?
- chrono::Duration::days(7))
.to_string();
client
.create_shard_key(
CreateShardKeyRequestBuilder::new(collection_name)
.request(CreateShardKeyBuilder::default().shard_key(today.to_string())),
)
.await?;
client
.delete_shard_key(
DeleteShardKeyRequestBuilder::new(collection_name)
.key(shard_key::Key::Keyword(oldest_shard_key)),
)
.await?;
String today = "2026-04-08";
String oldestShardKey = LocalDate.parse(today).minusDays(7).toString();
client.createShardKeyAsync(
CreateShardKeyRequest.newBuilder()
.setCollectionName(collectionName)
.setRequest(CreateShardKey.newBuilder()
.setShardKey(shardKey(today))
.build())
.build()
).get();
client.deleteShardKeyAsync(
DeleteShardKeyRequest.newBuilder()
.setCollectionName(collectionName)
.setRequest(DeleteShardKey.newBuilder()
.setShardKey(shardKey(oldestShardKey))
.build())
.build()
).get();
string today = "2026-04-08";
string oldestShardKey = DateOnly.ParseExact(today, "yyyy-MM-dd")
.AddDays(-7)
.ToString("yyyy-MM-dd");
await client.CreateShardKeyAsync(
collectionName,
new CreateShardKey { ShardKey = new ShardKey { Keyword = today } }
);
await client.DeleteShardKeyAsync(
collectionName,
new DeleteShardKey { ShardKey = new ShardKey { Keyword = oldestShardKey } }
);
today := "2026-04-08"
t, _ := time.Parse("2006-01-02", today)
oldestShardKey := t.AddDate(0, 0, -7).Format("2006-01-02")
client.CreateShardKey(context.Background(), collectionName, &qdrant.CreateShardKey{
ShardKey: qdrant.NewShardKey(today),
})
client.DeleteShardKey(context.Background(), collectionName, &qdrant.DeleteShardKey{
ShardKey: qdrant.NewShardKey(oldestShardKey),
})
새 데이터 수집
새 데이터를 수집할 때는 데이터가 올바른 샤드로 가도록 shard_key_selector를 오늘 날짜로 설정해요:
client.upsert(
collection_name=collection_name,
points=[PointStruct(
id=uuid.uuid4().hex,
payload={"text": "The best way to start a Wednesday is with a cup of coffee", "datetime": "2026-04-08T07:57:47"},
vector={
"dense_vector": Document(text="The best way to start a Wednesday is with a cup of coffee", model=dense_model)
})],
shard_key_selector=today
)
await client.upsert(collectionName, {
points: [
{
id: crypto.randomUUID(),
vector: {
dense_vector: {
text: "The best way to start a Wednesday is with a cup of coffee",
model: denseModel,
},
},
payload: {
text: "The best way to start a Wednesday is with a cup of coffee",
datetime: "2026-04-08T07:57:47",
},
},
],
shard_key: today,
});
client
.upsert_points(
UpsertPointsBuilder::new(
collection_name,
vec![PointStruct::new(
uuid::Uuid::new_v4().to_string(),
HashMap::from([(
"dense_vector".to_string(),
DocumentBuilder::new(
"The best way to start a Wednesday is with a cup of coffee",
dense_model,
)
.build(),
)]),
[
("text", "The best way to start a Wednesday is with a cup of coffee".into()),
("datetime", "2026-04-08T07:57:47".into()),
],
)],
)
.shard_key_selector(today.to_string()),
)
.await?;
client.upsertAsync(
UpsertPoints.newBuilder()
.setCollectionName(collectionName)
.addAllPoints(List.of(
PointStruct.newBuilder()
.setId(id(UUID.randomUUID()))
.setVectors(namedVectors(Map.of(
"dense_vector",
vector(Document.newBuilder()
.setText("The best way to start a Wednesday is with a cup of coffee")
.setModel(denseModel)
.build()))))
.putAllPayload(Map.of(
"text", value("The best way to start a Wednesday is with a cup of coffee"),
"datetime", value("2026-04-08T07:57:47")))
.build()))
.setShardKeySelector(shardKeySelector(today))
.build()
).get();
await client.UpsertAsync(
collectionName: collectionName,
points: new List<PointStruct>
{
new()
{
Id = Guid.NewGuid(),
Vectors = new Dictionary<string, Vector>
{
["dense_vector"] = new Document
{
Text = "The best way to start a Wednesday is with a cup of coffee",
Model = denseModel
}
},
Payload =
{
["text"] = "The best way to start a Wednesday is with a cup of coffee",
["datetime"] = "2026-04-08T07:57:47"
}
}
},
shardKeySelector: new ShardKeySelector
{
ShardKeys = { new List<ShardKey> { today } }
}
);
client.Upsert(context.Background(), &qdrant.UpsertPoints{
CollectionName: collectionName,
Points: []*qdrant.PointStruct{
{
Id: qdrant.NewID(uuid.New().String()),
Vectors: qdrant.NewVectorsMap(map[string]*qdrant.Vector{
"dense_vector": qdrant.NewVectorDocument(&qdrant.Document{
Text: "The best way to start a Wednesday is with a cup of coffee",
Model: denseModel,
}),
}),
Payload: qdrant.NewValueMap(map[string]any{
"text": "The best way to start a Wednesday is with a cup of coffee",
"datetime": "2026-04-08T07:57:47",
}),
},
},
ShardKeySelector: &qdrant.ShardKeySelector{
ShardKeys: []*qdrant.ShardKey{qdrant.NewShardKey(today)},
},
})
결론
시간 기반 샤딩은 Qdrant에서 크고 시계열적인 데이터셋을 관리하는 강력한 기법이에요. 타임스탬프를 기준으로 데이터를 서로 다른 샤드로 라우팅하면 최근 데이터를 효율적으로 저장하고 질의하면서, 성능에 영향을 주지 않고 오래된 데이터를 쉽게 정리할 수 있어요. 이 접근 방식은 데이터 관련성이 시간이 지나며 낮아지는 소셜 미디어 분석 같은 유스 케이스에 이상적이에요.