apache / apache/hudi

[SUPPORT]the compaction of the MOR hudi table keeps the old values

Open
#7,897 6 comments 0 reactions 1 assignee Claimed by @danny0405 View on GitHub
area:index area:table-service engine:spark
Dominant language
Java
Stars
6.2k
Forks
2.5k
Avg merge
2d 8h
Merged PRs (30d)
111

Description

I am having a glue job in which I write to hudi table, and I write it as MOR here's the config:

conf = {
'className': 'org.apache.hudi',
'hoodie.table.name': hudi_table_name,
'hoodie.datasource.write.operation': 'upsert',
'hoodie.datasource.write.table.type': 'MERGE_ON_READ',
'hoodie.datasource.write.precombine.field': 'timestamp',
'hoodie.datasource.write.recordkey.field': 'user_id',
#'hoodie.datasource.write.partitionpath.field': 'year:SIMPLE,month:SIMPLE,day:SIMPLE',
#'hoodie.datasource.write.keygenerator.class': 'org.apache.hudi.keygen.CustomKeyGenerator',
#'hoodie.deltastreamer.keygen.timebased.timestamp.type': 'DATE_STRING',
#'hoodie.deltastreamer.keygen.timebased.input.dateformat': 'yyyy-mm-dd',
#'hoodie.deltastreamer.keygen.timebased.output.dateformat': 'yyyy/MM/dd'
}

hudiGlueConfig = {
'hoodie.datasource.hive_sync.enable': 'true',
'hoodie.datasource.hive_sync.sync_as_datasource': 'false',
'hoodie.datasource.hive_sync.database': database_name,
'hoodie.datasource.hive_sync.table': hudi_table_name,
'hoodie.datasource.hive_sync.use_jdbc': 'false',
#'hoodie.datasource.write.hive_style_partitioning': 'false',
#'hoodie.datasource.hive_sync.partition_extractor_class': 'org.apache.hudi.hive.MultiPartKeysValueExtractor',
#'hoodie.datasource.hive_sync.partition_fields': 'year,month,day'
}

config_={**conf, **hudiGlueConfig}

I noticed that for each new record I append I had parquet file,so, first parquet has the first record, then when i insert new row a second parquet file created with both records, and when I insert for the third time a third parquet file is created with the 3 rows and when I update any of them I have a log file contains the update, and after number of appends the parqeut files compacted into one parquet file(the newest parquet file is kept (which has the three records appended) however , the other two parquet files are removed. However, this file contains the old values of initially added records, not the updated ones, any clue what I might be doing wrong?

the rt view reflects the correct data, the ro doesn't.

I am writing it as:
glueContext.forEachBatch( frame=data_frame_DataSource0, batch_function=processBatch, options={ "windowSize": window_size, "checkpointLocation": s3_path_spark } )

glueContext.write_dynamic_frame.from_options(
frame=DynamicFrame.fromDF(df, glueContext, "df"),
connection_type="custom.spark",
connection_options=config_
)

is it expected for each time I insert new record, a parquet file is created with accumulative records that were added? shouldn't it be reflected in the logfiles? in my case only the updates in existing records reflected in delta log files, but it's never written to the parquet files, even when I reach the number of commits, and the compaction is happening by deleting the oldest files and keeping only the last one created with the 3 rows inserted before separately but now they are in the same parquet?

case:

when I insert a row (id=3,name=mg) to the db:

spark streaming job creates a parquet file in the s3 path for hudi that contains (id=3,name=mg) --> file1.parquet

and the record is reflected in both rt, and ro tables

then when I add a new row (id=4,name=sa) :

spark streaming job creates a parquet file in the s3 path for hudi that contains both records --->file2.paruet
(id=3,name=mg)
(id=4,name=sa)

and both records is reflected in both rt, and ro tables

then when I add a new row (id=5,name=john) :

spark streaming job creates a parquet file in the s3 path for hudi that contains the three records --->file3.paruet
(id=3,name=mg)
(id=4,name=sa)
(id=5,name=john)

and the three records is reflected in both rt, and ro tables.

when I do multiple updates (say 9 updates) to the last record, I can see all these updates in the log files so my bucket contains 3 parquet files & 9 log files,

then after the 10th update where i changed the name to "joe", I can see 10 log files, and only 1 parquet file, the parquet file that is kept is the last one (file3.parquet) with the old values not the updates ones:
(id=3,name=mg)
(id=4,name=sa)
(id=5,name=john)

and file1.parquet &file2.parquet were delted.
rt table contained the right values (the three records and the last record has a value joe for the coloum name)
ro contained the values that's in the parquet
I
was expecting to find totally newly created parquet file that contains the values:
(id=3,name=mg)
(id=4,name=sa)
(id=5,name=joe)
and the ro to have the most updated ones.

Contributor guide

No contributing guide indexed for this repository

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.