Skip to content

Commit 0216742

Browse files
obdevfootkamiyuan-ljr
authored andcommitted
[CP] [fix] insert on dup error 5500
Co-authored-by: footka <672528926@qq.com> Co-authored-by: 884244693 <884244693@qq.com>
1 parent 8664da6 commit 0216742

6 files changed

Lines changed: 94 additions & 15 deletions

src/share/vector_index/ob_hybrid_vector_refresh_task.cpp

Lines changed: 14 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -425,7 +425,10 @@ int ObHybridVectorRefreshTask::get_embedded_table_column_ids(ObPluginVectorIndex
425425
}
426426
}
427427
task_ctx->part_key_num_ = part_key_column_ids.count();
428-
if (OB_SUCC(ret) && OB_FAIL(task_ctx->embedded_table_column_ids_.push_back(vector_column_id))) {
428+
if (OB_FAIL(ret)) {
429+
} else if (OB_FAIL(task_ctx->embedded_table_column_ids_.push_back(vector_column_id))) {
430+
LOG_WARN("failed to push vector column id.", K(ret));
431+
} else if (OB_FAIL(task_ctx->embedded_table_update_ids_.push_back(vector_column_id))) {
429432
LOG_WARN("failed to push vector column id.", K(ret));
430433
}
431434

@@ -902,7 +905,7 @@ int ObHybridVectorRefreshTask::after_embedding(ObPluginVectorIndexAdaptor &adapt
902905
ObArray<float*> output_vector;
903906
float *vector_buf = nullptr;
904907
ObHybridVectorRefreshTaskCtx *task_ctx = static_cast<ObHybridVectorRefreshTaskCtx *>(get_task_ctx());
905-
storage::ObValueRowIterator embedded_iter;
908+
ObVecIndexATaskUpdIterator embedded_iter;
906909
storage::ObValueRowIterator index_id_iter;
907910
storage::ObValueRowIterator delta_delete_iter;
908911
embedded_iter.init();
@@ -946,13 +949,15 @@ int ObHybridVectorRefreshTask::after_embedding(ObPluginVectorIndexAdaptor &adapt
946949
ret = OB_ALLOCATE_MEMORY_FAILED;
947950
LOG_WARN("failed to alloc mem.", K(ret), K(dim));
948951
} else {
949-
HEAP_VARS_3((blocksstable::ObDatumRow, datum_row, tenant_id_), (storage::ObTableScanParam, vid_rowkey_scan_param), (schema::ObTableParam, vid_rowkey_table_param, allocator_)) {
952+
HEAP_VARS_4((blocksstable::ObDatumRow, datum_row, tenant_id_), (blocksstable::ObDatumRow, new_row, tenant_id_), (storage::ObTableScanParam, vid_rowkey_scan_param), (schema::ObTableParam, vid_rowkey_table_param, allocator_)) {
950953
ObArenaAllocator scan_allocator("VecEmbedding", OB_MALLOC_NORMAL_BLOCK_SIZE, MTL_ID());
951954
common::ObNewRowIterator *vid_rowkey_iter = nullptr;
952955
ObTableScanIterator *table_scan_iter = nullptr;
953956
int64_t loop_cnt = 0;
954957
if (OB_FAIL(datum_row.init(task_ctx->embedded_table_column_ids_.count()))) {
955958
LOG_WARN("fail to init datum row", K(ret), K(task_ctx->embedded_table_column_ids_), K(datum_row));
959+
} else if (OB_FAIL(new_row.init(task_ctx->embedded_table_column_ids_.count()))) {
960+
LOG_WARN("fail to init datum row", K(ret), K(task_ctx->embedded_table_column_ids_), K(new_row));
956961
} else if (adaptor.get_is_need_vid() && OB_FAIL(ObPluginVectorIndexUtils::read_local_tablet(ls_id_,
957962
&adaptor,
958963
ctx_->task_status_.target_scn_,
@@ -970,7 +975,7 @@ int ObHybridVectorRefreshTask::after_embedding(ObPluginVectorIndexAdaptor &adapt
970975
for (int64_t i = 0; i < dim; i++) {
971976
vector_buf[i] = output_vector.at(row_id)[i];
972977
}
973-
datum_row.storage_datums_[task_ctx->embedded_table_column_ids_.count() - 1].set_string(reinterpret_cast<char *>(vector_buf), dim * sizeof(float));
978+
datum_row.storage_datums_[task_ctx->embedded_table_column_ids_.count() - 1].set_null();
974979
for (int i = task_ctx->embedded_table_column_ids_.count() - 2; i > task_ctx->embedded_table_column_ids_.count() - 2 - task_ctx->part_key_num_; i--) {
975980
datum_row.storage_datums_[i].set_null();
976981
}
@@ -1007,11 +1012,13 @@ int ObHybridVectorRefreshTask::after_embedding(ObPluginVectorIndexAdaptor &adapt
10071012
}
10081013
}
10091014
if (OB_SUCC(ret)) {
1010-
if (OB_FAIL(embedded_iter.add_row(datum_row))) {
1015+
if (OB_FAIL(new_row.deep_copy(datum_row, allocator_))) {
1016+
LOG_WARN("failed to copy row", K(ret), K(datum_row));
1017+
} else if (FALSE_IT(new_row.storage_datums_[task_ctx->embedded_table_column_ids_.count() - 1].set_string(reinterpret_cast<char *>(vector_buf), dim * sizeof(float)))) {
1018+
} else if (OB_FAIL(embedded_iter.add_row(datum_row, new_row))) {
10111019
LOG_WARN("failed to add row to index id iter", K(ret));
10121020
}
10131021
datum_row.reuse();
1014-
10151022
}
10161023

10171024
CHECK_TASK_CANCELLED_IN_PROCESS(ret, loop_cnt, ctx_);
@@ -1039,7 +1046,7 @@ int ObHybridVectorRefreshTask::after_embedding(ObPluginVectorIndexAdaptor &adapt
10391046
LOG_WARN("unexpected error", K(ret), KPC(task_ctx), K(oas));
10401047
} else if (OB_FAIL(init_dml_param(adaptor.get_embedded_table_id(), dml_param, table_dml_param, task_ctx->embedded_table_column_ids_, tx_desc, snapshot, store_ctx_guard))) {
10411048
LOG_WARN("failed to init dml param", K(ret), K(dml_param), K(table_dml_param));
1042-
} else if (OB_FAIL(oas->insert_rows(ls_id_, adaptor.get_embedded_tablet_id(), *tx_desc, dml_param, task_ctx->embedded_table_column_ids_, &embedded_iter, affected_rows))) {
1049+
} else if (OB_FAIL(oas->update_rows(ls_id_, adaptor.get_embedded_tablet_id(), *tx_desc, dml_param, task_ctx->embedded_table_column_ids_, task_ctx->embedded_table_update_ids_, &embedded_iter, affected_rows))) {
10431050
LOG_WARN("failed to insert rows to embedded table", K(ret), K(adaptor.get_embedded_tablet_id()));
10441051
}
10451052
store_ctx_guard.reset();

src/share/vector_index/ob_hybrid_vector_refresh_task.h

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -72,6 +72,7 @@ struct ObHybridVectorRefreshTaskCtx : public ObVecIndexAsyncTaskCtx
7272
embedding_task_(nullptr),
7373
index_id_column_ids_(),
7474
embedded_table_column_ids_(),
75+
embedded_table_update_ids_(),
7576
ai_service_(),
7677
endpoint_(),
7778
adp_guard_(),
@@ -97,6 +98,7 @@ struct ObHybridVectorRefreshTaskCtx : public ObVecIndexAsyncTaskCtx
9798
ObEmbeddingTask *embedding_task_;
9899
ObSEArray<uint64_t, 4> index_id_column_ids_;
99100
ObSEArray<uint64_t, 4> embedded_table_column_ids_;
101+
ObSEArray<uint64_t, 4> embedded_table_update_ids_;
100102
omt::ObAiServiceGuard ai_service_;
101103
const ObAiModelEndpointInfo *endpoint_;
102104
ObPluginVectorIndexAdapterGuard adp_guard_;

src/share/vector_index/ob_vector_index_async_task_util.cpp

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1184,6 +1184,47 @@ int ObVecIndexIAsyncTask::init(
11841184
return ret;
11851185
}
11861186

1187+
int ObVecIndexATaskUpdIterator::init() {
1188+
int ret = OB_SUCCESS;
1189+
if (OB_FAIL(old_row_.init())) {
1190+
LOG_WARN("fail to init old rows iter", K(ret));
1191+
} else if (OB_FAIL(new_row_.init())) {
1192+
LOG_WARN("fail to init new rows iter", K(ret));
1193+
}
1194+
return ret;
1195+
}
1196+
1197+
int ObVecIndexATaskUpdIterator::add_row(blocksstable::ObDatumRow &old_datum_row, blocksstable::ObDatumRow &new_datum_row) {
1198+
int ret = OB_SUCCESS;
1199+
if (OB_FAIL(old_row_.add_row(old_datum_row))) {
1200+
LOG_WARN("failed to add row to iter", K(ret));
1201+
} else if (OB_FAIL(new_row_.add_row(new_datum_row))) {
1202+
LOG_WARN("fail to init new rows iter", K(ret));
1203+
}
1204+
return ret;
1205+
}
1206+
1207+
int ObVecIndexATaskUpdIterator::get_next_row(blocksstable::ObDatumRow *&row)
1208+
{
1209+
int ret = OB_SUCCESS;
1210+
if (!got_old_row_) {
1211+
got_old_row_ = true;
1212+
if (OB_FAIL(old_row_.get_next_row(row))) {
1213+
if (OB_ITER_END != ret) {
1214+
LOG_WARN("fail to get next old row", K(ret));
1215+
}
1216+
}
1217+
} else {
1218+
got_old_row_ = false;
1219+
if (OB_FAIL(new_row_.get_next_row(row))) {
1220+
if (OB_ITER_END != ret) {
1221+
LOG_WARN("fail to get next new row", K(ret));
1222+
}
1223+
}
1224+
}
1225+
return ret;
1226+
}
1227+
11871228
/**************************** ObVecIndexAsyncTask ******************************/
11881229
int ObVecIndexAsyncTask::do_work()
11891230
{

src/share/vector_index/ob_vector_index_async_task_util.h

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@
2525
#include "lib/thread/thread_mgr_interface.h"
2626
#include "storage/access/ob_dml_param.h"
2727
#include "storage/tx/ob_trans_define_v4.h"
28+
#include "storage/ob_value_row_iterator.h"
2829

2930
namespace oceanbase
3031
{
@@ -316,6 +317,36 @@ class ObVecIndexIAsyncTask
316317
DISALLOW_COPY_AND_ASSIGN(ObVecIndexIAsyncTask);
317318
};
318319

320+
class ObVecIndexATaskUpdIterator : public blocksstable::ObDatumRowIterator
321+
{
322+
public:
323+
ObVecIndexATaskUpdIterator()
324+
: got_old_row_(false),
325+
is_iter_end_(false)
326+
{}
327+
328+
virtual ~ObVecIndexATaskUpdIterator() {
329+
old_row_.reset();
330+
new_row_.reset();
331+
}
332+
333+
int init();
334+
int add_row(blocksstable::ObDatumRow &old_datum_row, blocksstable::ObDatumRow &new_datum_row);
335+
336+
virtual int get_next_row(blocksstable::ObDatumRow *&row) override;
337+
virtual void reset() override {}
338+
339+
private:
340+
// disallow copy
341+
DISALLOW_COPY_AND_ASSIGN(ObVecIndexATaskUpdIterator);
342+
343+
private:
344+
storage::ObValueRowIterator old_row_;
345+
storage::ObValueRowIterator new_row_;
346+
bool got_old_row_;
347+
bool is_iter_end_;
348+
};
349+
319350
class ObVecIndexAsyncTask : public ObVecIndexIAsyncTask
320351
{
321352
public:

src/sql/das/ob_das_dml_vec_iter.cpp

Lines changed: 5 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -686,16 +686,14 @@ int ObEmbeddedVecDMLIterator::generate_domain_rows(const ObChunkDatumStore::Stor
686686
bool is_sync_interval = false;
687687
if (OB_FAIL(check_sync_interval(is_sync_interval))) {
688688
LOG_WARN("fail to check sync interval", K(ret));
689-
} else if (is_sync_interval && OB_FAIL(generate_embedded_vec_row(store_row))) {
689+
} else if (OB_FAIL(generate_embedded_vec_row(store_row, is_sync_interval))) {
690690
LOG_WARN("failed to generate embedded vec row", K(ret));
691-
} else if (!is_sync_interval) {
692-
ret = OB_ITER_END;
693691
}
694692
}
695693
return ret;
696694
}
697695

698-
int ObEmbeddedVecDMLIterator::generate_embedded_vec_row(const ObChunkDatumStore::StoredRow *store_row)
696+
int ObEmbeddedVecDMLIterator::generate_embedded_vec_row(const ObChunkDatumStore::StoredRow *store_row, bool is_sync)
699697
{
700698
int ret = OB_SUCCESS;
701699

@@ -729,14 +727,14 @@ int ObEmbeddedVecDMLIterator::generate_embedded_vec_row(const ObChunkDatumStore:
729727
int64_t vid = OB_INVALID_ID;
730728
if (OB_FAIL(get_vid(store_row, vid_idx, vid))) {
731729
LOG_WARN("failed to get vid", K(ret));
730+
} else if (!is_sync) {
731+
obj_arr[vid_idx].set_int(vid);
732+
obj_arr[embedded_vec_idx].set_null();
732733
} else if (OB_FAIL(get_chunk_data(store_row, embedded_vec_idx, chunk))) {
733734
LOG_WARN("failed to project chunk columns for embedding", K(ret));
734735
} else if (!is_old_row_ && chunk.empty()) {
735736
obj_arr[vid_idx].set_int(vid);
736737
obj_arr[embedded_vec_idx].set_null();
737-
if (OB_FAIL(rows_.push_back(row))) {
738-
LOG_WARN("fail to push back row", K(ret));
739-
}
740738
} else {
741739
obj_arr[vid_idx].set_int(vid);
742740
ObString embedded_vector;

src/sql/das/ob_das_dml_vec_iter.h

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -159,7 +159,7 @@ class ObEmbeddedVecDMLIterator final : public ObDomainDMLIterator
159159
protected:
160160
int generate_domain_rows(const ObChunkDatumStore::StoredRow *store_row) override;
161161
private:
162-
int generate_embedded_vec_row(const ObChunkDatumStore::StoredRow *store_row);
162+
int generate_embedded_vec_row(const ObChunkDatumStore::StoredRow *store_row, bool is_sync);
163163
int get_embedded_vec_column_idxs(int64_t &vid_idx, int64_t &embedded_vec_idx);
164164
int get_vid(const ObChunkDatumStore::StoredRow *store_row, const int64_t vid_idx, int64_t &vid);
165165
int get_chunk_data(const ObChunkDatumStore::StoredRow *store_row, const int64_t embedded_vec_idx, ObString &chunk);

0 commit comments

Comments
 (0)