Skip to content

Commit 2881ded

Browse files
authored
Temp commit (#1)
1 parent 9e69f9a commit 2881ded

38 files changed

Lines changed: 2872 additions & 2433 deletions

DEPS.bzl

Lines changed: 943 additions & 1307 deletions
Large diffs are not rendered by default.

br/pkg/storage/gcs.go

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -594,8 +594,6 @@ type gcsObjectReader struct {
594594

595595
prefetchSize int
596596
// reader context used for implement `io.Seek`
597-
// currently, lightning depends on package `xitongsys/parquet-go` to read parquet file and it needs `io.Seeker`
598-
// See: https://github.com/xitongsys/parquet-go/blob/207a3cee75900b2b95213627409b7bac0f190bb3/source/source.go#L9-L10
599597
ctx context.Context
600598
}
601599

br/pkg/storage/s3.go

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1018,8 +1018,6 @@ type s3ObjectReader struct {
10181018
pos int64
10191019
rangeInfo RangeInfo
10201020
// reader context used for implement `io.Seek`
1021-
// currently, lightning depends on package `xitongsys/parquet-go` to read parquet file and it needs `io.Seeker`
1022-
// See: https://github.com/xitongsys/parquet-go/blob/207a3cee75900b2b95213627409b7bac0f190bb3/source/source.go#L9-L10
10231021
ctx context.Context
10241022
prefetchSize int
10251023
}

go.mod

Lines changed: 8 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ require (
1313
github.com/Masterminds/semver v1.5.0
1414
github.com/YangKeao/go-mysql-driver v0.0.0-20240627104025-dd5589458cfa
1515
github.com/aliyun/alibaba-cloud-sdk-go v1.61.1581
16+
github.com/apache/arrow-go/v18 v18.0.0
1617
github.com/apache/skywalking-eyes v0.4.0
1718
github.com/asaskevich/govalidator v0.0.0-20230301143203-a9d515a09cc2
1819
github.com/ashanbrown/forbidigo/v2 v2.1.0
@@ -118,8 +119,6 @@ require (
118119
github.com/uber/jaeger-client-go v2.22.1+incompatible
119120
github.com/vbauerster/mpb/v7 v7.5.3
120121
github.com/wangjohn/quickselect v0.0.0-20161129230411-ed8402a42d5f
121-
github.com/xitongsys/parquet-go v1.6.3-0.20240520233950-75e935fc3e17
122-
github.com/xitongsys/parquet-go-source v0.0.0-20200817004010-026bad9b25d0
123122
github.com/zyedidia/generic v1.2.1
124123
go.etcd.io/etcd/api/v3 v3.5.15
125124
go.etcd.io/etcd/client/pkg/v3 v3.5.15
@@ -156,24 +155,22 @@ require (
156155
require (
157156
codeberg.org/chavacava/garif v0.2.0 // indirect
158157
filippo.io/edwards25519 v1.1.0 // indirect
159-
github.com/andybalholm/brotli v1.0.5 // indirect
160-
github.com/apache/arrow/go/v12 v12.0.1 // indirect
158+
github.com/andybalholm/brotli v1.1.1 // indirect
161159
github.com/cockroachdb/errors v1.11.3 // indirect
162160
github.com/cockroachdb/fifo v0.0.0-20240606204812-0bbfbd93a7ce // indirect
163161
github.com/cockroachdb/tokenbucket v0.0.0-20230807174530-cc333fc44b06 // indirect
164162
github.com/getsentry/sentry-go v0.27.0 // indirect
165-
github.com/goccy/go-reflect v1.2.0 // indirect
166-
github.com/google/flatbuffers v2.0.8+incompatible // indirect
163+
github.com/google/flatbuffers v24.3.25+incompatible // indirect
167164
github.com/google/go-cmp v0.7.0 // indirect
168165
github.com/jinzhu/inflection v1.0.0 // indirect
169166
github.com/jinzhu/now v1.1.5 // indirect
170167
github.com/klauspost/asmfmt v1.3.2 // indirect
171-
github.com/klauspost/cpuid/v2 v2.2.7 // indirect
168+
github.com/klauspost/cpuid/v2 v2.2.9 // indirect
172169
github.com/ldez/grignotin v0.9.0 // indirect
173170
github.com/minio/asm2plan9s v0.0.0-20200509001527-cdd76441f9d8 // indirect
174171
github.com/minio/c2goasm v0.0.0-20190812172519-36a3d3bbc4f3 // indirect
175172
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
176-
github.com/pierrec/lz4/v4 v4.1.15 // indirect
173+
github.com/pierrec/lz4/v4 v4.1.21 // indirect
177174
github.com/qri-io/jsonpointer v0.1.1 // indirect
178175
github.com/segmentio/fasthash v1.0.3 // indirect
179176
github.com/tidwall/gjson v1.14.4 // indirect
@@ -197,7 +194,7 @@ require (
197194
github.com/Masterminds/sprig/v3 v3.2.2 // indirect
198195
github.com/VividCortex/ewma v1.2.0 // indirect
199196
github.com/acarl005/stripansi v0.0.0-20180116102854-5a71ef0e047d // indirect
200-
github.com/apache/thrift v0.16.0 // indirect
197+
github.com/apache/thrift v0.21.0 // indirect
201198
github.com/beorn7/perks v1.0.1 // indirect
202199
github.com/bmatcuk/doublestar/v2 v2.0.4 // indirect
203200
github.com/cenkalti/backoff/v4 v4.2.1 // indirect
@@ -220,7 +217,7 @@ require (
220217
github.com/go-logr/logr v1.4.1 // indirect
221218
github.com/go-logr/stdr v1.2.2 // indirect
222219
github.com/go-ole/go-ole v1.3.0 // indirect
223-
github.com/goccy/go-json v0.10.2 // indirect
220+
github.com/goccy/go-json v0.10.4 // indirect
224221
github.com/golang-jwt/jwt/v4 v4.5.2 // indirect
225222
github.com/golang-jwt/jwt/v5 v5.2.2 // indirect
226223
github.com/golang/glog v1.2.4 // indirect
@@ -330,6 +327,7 @@ require (
330327
)
331328

332329
replace (
330+
github.com/apache/arrow-go/v18 => github.com/joechenrh/arrow-go/v18 v18.0.0-20250905011811-90682f7df921
333331
github.com/go-ldap/ldap/v3 => github.com/YangKeao/ldap/v3 v3.4.5-0.20230421065457-369a3bab1117
334332
github.com/pingcap/tidb/pkg/parser => ./pkg/parser
335333

go.sum

Lines changed: 26 additions & 191 deletions
Large diffs are not rendered by default.

lightning/pkg/importer/BUILD.bazel

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -182,8 +182,6 @@ go_test(
182182
"@com_github_tikv_pd_client//:client",
183183
"@com_github_tikv_pd_client//http",
184184
"@com_github_tikv_pd_client//pkg/caller",
185-
"@com_github_xitongsys_parquet_go//writer",
186-
"@com_github_xitongsys_parquet_go_source//buffer",
187185
"@io_etcd_go_etcd_client_v3//:client",
188186
"@io_etcd_go_etcd_tests_v3//integration",
189187
"@org_uber_go_mock//gomock",

lightning/pkg/importer/chunk_process.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -107,7 +107,7 @@ func openParser(
107107
case mydump.SourceTypeSQL:
108108
parser = mydump.NewChunkParser(ctx, cfg.TiDB.SQLMode, reader, blockBufSize, ioWorkers)
109109
case mydump.SourceTypeParquet:
110-
parser, err = mydump.NewParquetParser(ctx, store, reader, chunk.FileMeta.Path)
110+
parser, err = mydump.NewParquetParser(ctx, store, reader, chunk.FileMeta.Path, chunk.FileMeta.ParquetMeta)
111111
if err != nil {
112112
return nil, err
113113
}

lightning/pkg/importer/get_pre_info.go

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -491,7 +491,10 @@ func (p *PreImportInfoGetterImpl) ReadFirstNRowsByFileMeta(ctx context.Context,
491491
case mydump.SourceTypeSQL:
492492
parser = mydump.NewChunkParser(ctx, p.cfg.TiDB.SQLMode, reader, blockBufSize, p.ioWorkers)
493493
case mydump.SourceTypeParquet:
494-
parser, err = mydump.NewParquetParser(ctx, p.srcStorage, reader, dataFileMeta.Path)
494+
parser, err = mydump.NewParquetParser(
495+
ctx, p.srcStorage, reader,
496+
dataFileMeta.Path, mydump.GetDefaultParquetMeta(),
497+
)
495498
if err != nil {
496499
return nil, nil, errors.Trace(err)
497500
}
@@ -661,7 +664,10 @@ func (p *PreImportInfoGetterImpl) sampleDataFromTable(
661664
case mydump.SourceTypeSQL:
662665
parser = mydump.NewChunkParser(ctx, p.cfg.TiDB.SQLMode, reader, blockBufSize, p.ioWorkers)
663666
case mydump.SourceTypeParquet:
664-
parser, err = mydump.NewParquetParser(ctx, p.srcStorage, reader, sampleFile.Path)
667+
parser, err = mydump.NewParquetParser(
668+
ctx, p.srcStorage, reader,
669+
sampleFile.Path, mydump.GetDefaultParquetMeta(),
670+
)
665671
if err != nil {
666672
return 0.0, false, errors.Trace(err)
667673
}

lightning/pkg/importer/get_pre_info_test.go

Lines changed: 18 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -20,13 +20,13 @@ import (
2020
"context"
2121
"database/sql"
2222
"fmt"
23-
"slices"
2423
"strings"
2524
"testing"
2625

2726
"github.com/DATA-DOG/go-sqlmock"
2827
mysql_sql_driver "github.com/go-sql-driver/mysql"
2928
"github.com/pingcap/errors"
29+
"github.com/pingcap/tidb/br/pkg/storage"
3030
"github.com/pingcap/tidb/lightning/pkg/importer/mock"
3131
ropts "github.com/pingcap/tidb/lightning/pkg/importer/opts"
3232
"github.com/pingcap/tidb/pkg/errno"
@@ -36,8 +36,6 @@ import (
3636
"github.com/pingcap/tidb/pkg/parser/ast"
3737
"github.com/pingcap/tidb/pkg/types"
3838
"github.com/stretchr/testify/require"
39-
pqt_buf_src "github.com/xitongsys/parquet-go-source/buffer"
40-
pqtwriter "github.com/xitongsys/parquet-go/writer"
4139
)
4240

4341
type colDef struct {
@@ -253,26 +251,24 @@ func TestGetPreInfoGetAllTableStructures(t *testing.T) {
253251
}
254252
}
255253

256-
func generateParquetData(t *testing.T) []byte {
257-
type parquetStruct struct {
258-
ID int64 `parquet:"name=id, type=INT64"`
259-
Name string `parquet:"name=name, type=BYTE_ARRAY"`
260-
}
261-
pf, err := pqt_buf_src.NewBufferFile(make([]byte, 0))
254+
func readParquetData(t *testing.T) []byte {
255+
s, err := storage.ParseBackend("./testdata", nil)
262256
require.NoError(t, err)
263-
pw, err := pqtwriter.NewParquetWriter(pf, new(parquetStruct), 4)
257+
258+
store, err := storage.NewWithDefaultOpt(context.Background(), s)
264259
require.NoError(t, err)
265-
for i := range 10 {
266-
require.NoError(t, pw.Write(parquetStruct{
267-
ID: int64(i + 1),
268-
Name: fmt.Sprintf("name_%d", i+1),
269-
}))
270-
}
271-
require.NoError(t, pw.WriteStop())
272-
require.NoError(t, pf.Close())
273-
bf, ok := pf.(pqt_buf_src.BufferFile)
274-
require.True(t, ok)
275-
return slices.Clone(bf.Bytes())
260+
defer store.Close()
261+
262+
reader, err := store.Open(context.Background(), "test.parquet", nil)
263+
require.NoError(t, err)
264+
defer reader.Close()
265+
266+
bs := make([]byte, 1024)
267+
l, err := reader.Read(bs)
268+
bs = bs[:l]
269+
require.NoError(t, err)
270+
271+
return bs
276272
}
277273

278274
func TestGetPreInfoReadFirstRow(t *testing.T) {
@@ -282,7 +278,6 @@ func TestGetPreInfoReadFirstRow(t *testing.T) {
282278
111,"aaa"
283279
222,"bbb"
284280
`)
285-
pqtData := generateParquetData(t)
286281
const testSQLData01 string = `INSERT INTO db01.tbl01 (ival, sval) VALUES (333, 'ccc');
287282
INSERT INTO db01.tbl01 (ival, sval) VALUES (444, 'ddd');`
288283
testDataInfos := []struct {
@@ -349,7 +344,7 @@ INSERT INTO db01.tbl01 (ival, sval) VALUES (444, 'ddd');`
349344
},
350345
{
351346
FileName: "/db01/tbl01/data.005.parquet",
352-
Data: pqtData,
347+
Data: readParquetData(t),
353348
FirstN: 3,
354349
ExpectFirstRowDatums: [][]types.Datum{
355350
{

lightning/pkg/importer/table_import.go

Lines changed: 0 additions & 44 deletions
Original file line numberDiff line numberDiff line change
@@ -29,7 +29,6 @@ import (
2929
dmysql "github.com/go-sql-driver/mysql"
3030
"github.com/pingcap/errors"
3131
"github.com/pingcap/failpoint"
32-
"github.com/pingcap/tidb/br/pkg/storage"
3332
"github.com/pingcap/tidb/br/pkg/version"
3433
"github.com/pingcap/tidb/lightning/pkg/web"
3534
"github.com/pingcap/tidb/pkg/errno"
@@ -788,13 +787,6 @@ ChunkLoop:
788787
break
789788
}
790789

791-
if chunk.FileMeta.Type == mydump.SourceTypeParquet {
792-
// TODO: use the compressed size of the chunk to conduct memory control
793-
if _, err = getChunkCompressedSizeForParquet(ctx, chunk, rc.store); err != nil {
794-
return nil, errors.Trace(err)
795-
}
796-
}
797-
798790
restoreWorker := rc.regionWorkers.Apply()
799791
wg.Add(1)
800792
go func(w *worker.Worker, cr *chunkProcessor) {
@@ -1201,42 +1193,6 @@ func (tr *TableImporter) postProcess(
12011193
return true, nil
12021194
}
12031195

1204-
func getChunkCompressedSizeForParquet(
1205-
ctx context.Context,
1206-
chunk *checkpoints.ChunkCheckpoint,
1207-
store storage.ExternalStorage,
1208-
) (int64, error) {
1209-
reader, err := mydump.OpenReader(ctx, &chunk.FileMeta, store, storage.DecompressConfig{
1210-
ZStdDecodeConcurrency: 1,
1211-
})
1212-
if err != nil {
1213-
return 0, errors.Trace(err)
1214-
}
1215-
parser, err := mydump.NewParquetParser(ctx, store, reader, chunk.FileMeta.Path)
1216-
if err != nil {
1217-
_ = reader.Close()
1218-
return 0, errors.Trace(err)
1219-
}
1220-
//nolint: errcheck
1221-
defer parser.Close()
1222-
err = parser.Reader.ReadFooter()
1223-
if err != nil {
1224-
return 0, errors.Trace(err)
1225-
}
1226-
rowGroups := parser.Reader.Footer.GetRowGroups()
1227-
var maxRowGroupSize int64
1228-
for _, rowGroup := range rowGroups {
1229-
var rowGroupSize int64
1230-
columnChunks := rowGroup.GetColumns()
1231-
for _, columnChunk := range columnChunks {
1232-
columnChunkSize := columnChunk.MetaData.GetTotalCompressedSize()
1233-
rowGroupSize += columnChunkSize
1234-
}
1235-
maxRowGroupSize = max(maxRowGroupSize, rowGroupSize)
1236-
}
1237-
return maxRowGroupSize, nil
1238-
}
1239-
12401196
func updateStatsMeta(ctx context.Context, db *sql.DB, tableID int64, count int) {
12411197
s := common.SQLWithRetry{
12421198
DB: db,

0 commit comments

Comments
 (0)