Skip to content

Commit 8ea1c3c

Browse files
committed
fixes
1 parent 2e899bb commit 8ea1c3c

6 files changed

Lines changed: 49 additions & 49 deletions

File tree

src/copy.c

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -574,7 +574,11 @@ TSCopyMultiInsertBufferCleanup(TSCopyMultiInsertInfo *miinfo, TSCopyMultiInsertB
574574
FreeTupleDesc(buffer->tupdesc);
575575
break;
576576
case TS_CIM_COMPRESSION:
577-
ts_cm_functions->compressor_free(buffer->compressor, buffer->bulk_writer);
577+
ts_cm_functions->compressor_close(buffer->compressor, buffer->bulk_writer);
578+
pfree(buffer->compressor);
579+
pfree(buffer->bulk_writer);
580+
buffer->compressor = NULL;
581+
buffer->bulk_writer = NULL;
578582
break;
579583
}
580584

src/cross_module_fn.h

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -135,7 +135,7 @@ typedef struct CrossModuleFunctions
135135
void (*compressor_add_slot)(RowCompressor *compressor, BulkWriter *bulk_writer,
136136
TupleTableSlot *slot);
137137
void (*compressor_flush)(RowCompressor *compressor, BulkWriter *bulk_writer);
138-
void (*compressor_free)(RowCompressor *compressor, BulkWriter *bulk_writer);
138+
void (*compressor_close)(RowCompressor *compressor, BulkWriter *bulk_writer);
139139
Chunk *(*compression_chunk_create)(Hypertable *ht, Chunk *src_chunk);
140140

141141
/* The compression functions below are not installed in SQL as part of create extension;

src/nodes/modify_hypertable_exec.c

Lines changed: 11 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -2485,10 +2485,13 @@ ExecModifyTable(CustomScanState *cs_node, PlanState *pstate)
24852485
/* Flush on chunk change */
24862486
if (ht_state->compressor && ht_state->compressor_relid != RelationGetRelid(ctr->cis->rel))
24872487
{
2488-
ts_cm_functions->compressor_flush(ht_state->compressor, ht_state->bulk_writer);
2489-
ts_cm_functions->compressor_free(ht_state->compressor, ht_state->bulk_writer);
2490-
ht_state->compressor = NULL;
2491-
ht_state->compressor_relid = InvalidOid;
2488+
ts_cm_functions->compressor_flush(ht_state->compressor, ht_state->bulk_writer);
2489+
ts_cm_functions->compressor_close(ht_state->compressor, ht_state->bulk_writer);
2490+
pfree(ht_state->compressor);
2491+
pfree(ht_state->bulk_writer);
2492+
ht_state->compressor = NULL;
2493+
ht_state->compressor_relid = InvalidOid;
2494+
ht_state->bulk_writer = NULL;
24922495
}
24932496

24942497
if (!ht_state->compressor)
@@ -2814,8 +2817,11 @@ ExecModifyTable(CustomScanState *cs_node, PlanState *pstate)
28142817
if (ht_state->compressor)
28152818
{
28162819
ts_cm_functions->compressor_flush(ht_state->compressor, ht_state->bulk_writer);
2817-
ts_cm_functions->compressor_free(ht_state->compressor, ht_state->bulk_writer);
2820+
ts_cm_functions->compressor_close(ht_state->compressor, ht_state->bulk_writer);
2821+
pfree(ht_state->compressor);
2822+
pfree(ht_state->bulk_writer);
28182823
ht_state->compressor = NULL;
2824+
ht_state->compressor_relid = InvalidOid;
28192825
ht_state->bulk_writer = NULL;
28202826
}
28212827

tsl/src/compression/compression.c

Lines changed: 30 additions & 36 deletions
Original file line numberDiff line numberDiff line change
@@ -157,6 +157,9 @@ static ArrayType *analyze_segmentby_candidates(ColumnAnalysis *candidates, int n
157157
static ArrayType *analyze_and_get_segmentby(CompressionSettings *settings,
158158
RowCompressor *compressor);
159159

160+
static void compressor_apply_segmentby_and_rebuild(RowCompressor *old_compressor,
161+
BulkWriter *old_bulk_writer);
162+
160163
/********************
161164
** compress_chunk **
162165
********************/
@@ -1063,6 +1066,9 @@ tsl_compressor_add_slot(RowCompressor *compressor, BulkWriter *bulk_writer, Tupl
10631066
void
10641067
tsl_compressor_flush(RowCompressor *compressor, BulkWriter *bulk_writer)
10651068
{
1069+
fprintf(stderr, "flush compressor at %p\n", compressor);
1070+
mybt();
1071+
10661072
if (compressor->sort_state)
10671073
{
10681074
if (compressor->tuples_to_sort)
@@ -1071,15 +1077,7 @@ tsl_compressor_flush(RowCompressor *compressor, BulkWriter *bulk_writer)
10711077

10721078
if (compressor->needs_analyze_segmentby)
10731079
{
1074-
RowCompressor *new_compressor = NULL;
1075-
BulkWriter *new_bulk_writer = NULL;
1076-
tsl_compressor_apply_segmentby_and_rebuild(compressor,
1077-
bulk_writer,
1078-
&new_compressor,
1079-
&new_bulk_writer);
1080-
*compressor = *new_compressor;
1081-
compressor->needs_analyze_segmentby = false;
1082-
*bulk_writer = *new_bulk_writer;
1080+
compressor_apply_segmentby_and_rebuild(compressor, bulk_writer);
10831081
}
10841082

10851083
TupleTableSlot *slot = MakeTupleTableSlot(compressor->in_desc, &TTSOpsMinimalTuple);
@@ -1116,8 +1114,11 @@ tsl_compressor_flush(RowCompressor *compressor, BulkWriter *bulk_writer)
11161114
}
11171115

11181116
void
1119-
tsl_compressor_free(RowCompressor *compressor, BulkWriter *bulk_writer)
1117+
tsl_compressor_close(RowCompressor *compressor, BulkWriter *bulk_writer)
11201118
{
1119+
fprintf(stderr, "free compressor at %p\n", compressor);
1120+
mybt();
1121+
11211122
if (compressor->sort_state)
11221123
{
11231124
tuplesort_end(compressor->sort_state);
@@ -1139,14 +1140,10 @@ tsl_compressor_free(RowCompressor *compressor, BulkWriter *bulk_writer)
11391140
* Determine the segmentby column from tuples in Tuplesortstate in the RowCompressor,
11401141
* then rebuild the compressed chunk and compressor to use it.
11411142
*/
1142-
void
1143-
tsl_compressor_apply_segmentby_and_rebuild(RowCompressor *old_compressor,
1144-
BulkWriter *old_bulk_writer,
1145-
RowCompressor **output_compressor,
1146-
BulkWriter **output_bulk_writer)
1143+
static void
1144+
compressor_apply_segmentby_and_rebuild(RowCompressor *old_compressor, BulkWriter *old_bulk_writer)
11471145
{
1148-
*output_bulk_writer = old_bulk_writer;
1149-
*output_compressor = old_compressor;
1146+
old_compressor->needs_analyze_segmentby = false;
11501147
if (old_compressor->sort_state == NULL || old_compressor->tuples_to_sort == 0)
11511148
{
11521149
return;
@@ -1225,19 +1222,20 @@ tsl_compressor_apply_segmentby_and_rebuild(RowCompressor *old_compressor,
12251222
/* Initialize the new bulk writer and compressor against the new compressed relation */
12261223
Relation out_rel = table_open(new_compressed_chunk->table_id, RowExclusiveLock);
12271224

1228-
BulkWriter *new_bulk_writer = bulk_writer_alloc(out_rel, /* insert_options = */ 0);
1225+
BulkWriter new_bulk_writer = bulk_writer_build(out_rel, /* insert_options = */ 0);
12291226

1230-
RowCompressor *new_compressor = palloc0(sizeof(RowCompressor));
1231-
row_compressor_init(new_compressor,
1227+
RowCompressor new_compressor;
1228+
row_compressor_init(&new_compressor,
12321229
settings,
12331230
RelationGetDescr(in_rel),
12341231
RelationGetDescr(out_rel));
12351232

1236-
MemoryContext old_context = MemoryContextSwitchTo(new_compressor->row_compressor_context);
1237-
new_compressor->sort_state = compression_create_tuplesort_state(settings, in_rel, false);
1233+
MemoryContext old_context = MemoryContextSwitchTo(new_compressor.row_compressor_context);
1234+
new_compressor.sort_state = compression_create_tuplesort_state(settings, in_rel, false);
12381235
MemoryContextSwitchTo(old_context);
12391236

1240-
new_compressor->tuple_sort_limit = old_compressor->tuple_sort_limit;
1237+
new_compressor.tuple_sort_limit = old_compressor->tuple_sort_limit;
1238+
new_compressor.needs_analyze_segmentby = false;
12411239

12421240
/* Transfer from old sort state into the new one with segmentby settings */
12431241
TupleTableSlot *slot = MakeTupleTableSlot(old_compressor->in_desc, &TTSOpsMinimalTuple);
@@ -1247,22 +1245,22 @@ tsl_compressor_apply_segmentby_and_rebuild(RowCompressor *old_compressor,
12471245
slot,
12481246
NULL /*=abbrev*/))
12491247
{
1250-
tuplesort_puttupleslot(new_compressor->sort_state, slot);
1251-
new_compressor->tuples_to_sort++;
1248+
tuplesort_puttupleslot(new_compressor.sort_state, slot);
1249+
new_compressor.tuples_to_sort++;
12521250
}
12531251

12541252
ExecDropSingleTupleTableSlot(slot);
12551253

1256-
tsl_compressor_free(old_compressor, old_bulk_writer);
1254+
tsl_compressor_close(old_compressor, old_bulk_writer);
12571255

12581256
ts_chunk_drop(old_compressed_chunk, DROP_RESTRICT, -1);
12591257

1260-
tuplesort_performsort(new_compressor->sort_state);
1258+
tuplesort_performsort(new_compressor.sort_state);
12611259

12621260
table_close(in_rel, NoLock);
12631261

1264-
*output_bulk_writer = new_bulk_writer;
1265-
*output_compressor = new_compressor;
1262+
*old_bulk_writer = new_bulk_writer;
1263+
*old_compressor = new_compressor;
12661264
}
12671265

12681266
/*
@@ -1278,6 +1276,8 @@ tsl_compressor_init(Relation in_rel, BulkWriter **bulk_writer, bool sort, int so
12781276
CompressionSettings *settings = ts_compression_settings_get(in_rel->rd_id);
12791277
Relation out_rel = table_open(settings->fd.compress_relid, RowExclusiveLock);
12801278
RowCompressor *compressor = palloc0(sizeof(RowCompressor));
1279+
fprintf(stderr, "alloc compressor at %p\n", compressor);
1280+
mybt();
12811281
row_compressor_init(compressor, settings, RelationGetDescr(in_rel), RelationGetDescr(out_rel));
12821282

12831283
*bulk_writer = bulk_writer_alloc(out_rel, 0);
@@ -1999,13 +1999,7 @@ BulkWriter *
19991999
bulk_writer_alloc(Relation out_rel, int insert_options)
20002000
{
20012001
BulkWriter *writer = palloc(sizeof(BulkWriter));
2002-
writer->out_rel = out_rel;
2003-
writer->indexstate = CatalogOpenIndexes(out_rel);
2004-
writer->mycid = GetCurrentCommandId(true);
2005-
writer->bistate = GetBulkInsertState();
2006-
writer->estate = CreateExecutorState();
2007-
writer->insert_options = insert_options;
2008-
2002+
*writer = bulk_writer_build(out_rel, insert_options);
20092003
return writer;
20102004
}
20112005

tsl/src/compression/compression.h

Lines changed: 1 addition & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -406,16 +406,12 @@ extern void row_compressor_init(RowCompressor *row_compressor, const Compression
406406

407407
extern RowCompressor *tsl_compressor_init(Relation in_rel, BulkWriter **bulk_writer, bool sort,
408408
int tuple_sort_limit, bool created_compressed_chunk);
409-
extern void tsl_compressor_apply_segmentby_and_rebuild(RowCompressor *old_compressor,
410-
BulkWriter *old_bulk_writer,
411-
RowCompressor **output_compressor,
412-
BulkWriter **output_bulk_writer);
413409
extern void tsl_compressor_set_invalidation(RowCompressor *compressor, Hypertable *ht,
414410
Oid chunk_relid);
415411
extern void tsl_compressor_add_slot(RowCompressor *compressor, BulkWriter *bulk_writer,
416412
TupleTableSlot *slot);
417413
extern void tsl_compressor_flush(RowCompressor *compressor, BulkWriter *bulk_writer);
418-
extern void tsl_compressor_free(RowCompressor *compressor, BulkWriter *bulk_writer);
414+
extern void tsl_compressor_close(RowCompressor *compressor, BulkWriter *bulk_writer);
419415

420416
extern void row_compressor_reset(RowCompressor *row_compressor);
421417
extern void row_compressor_close(RowCompressor *row_compressor);

tsl/src/init.c

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -177,7 +177,7 @@ CrossModuleFunctions tsl_cm_functions = {
177177
.compressor_set_invalidation = tsl_compressor_set_invalidation,
178178
.compressor_add_slot = tsl_compressor_add_slot,
179179
.compressor_flush = tsl_compressor_flush,
180-
.compressor_free = tsl_compressor_free,
180+
.compressor_close = tsl_compressor_close,
181181
.compression_chunk_create = tsl_compression_chunk_create,
182182
.show_chunk = chunk_show,
183183
.create_compressed_chunk = tsl_create_compressed_chunk,

0 commit comments

Comments
 (0)