diff --git a/README.md b/README.md index cf2dcc8..61ea505 100644 --- a/README.md +++ b/README.md @@ -55,5 +55,5 @@ For more details about parquet format , see [Parquet File Format](https://parque - If the S3 archive job fails, it is resumable , but there might be some data duplication in the destination due to the nature of chunking and the checkpointing mechanism. Future versions may address this issue. But until then you can refer to this example Glue(spark) job to remove duplicates from parquet files and insert to iceberg table. [Remove Duplicates from Parquet Files](examples/remove_duplicates_from_parquet_files.py) ## Dependencies -* Relies on [mysql-client-driver](https://github.com/go-sql-driver/mysql) for connecting to MySQL databases. +* Relies on [block/mysql](https://github.com/block/mysql), Block's fork of [go-sql-driver/mysql](https://github.com/go-sql-driver/mysql), for connecting to MySQL databases. It registers under the driver name `block-mysql`. * Relies on [Spirit](https://github.com/block/spirit) for chunking the data that needs to be archived and also for many other utilities related to db connection, loading table metadata etc. \ No newline at end of file diff --git a/go.mod b/go.mod index bac2497..9460b10 100644 --- a/go.mod +++ b/go.mod @@ -8,17 +8,16 @@ require ( github.com/aws/aws-sdk-go-v2 v1.33.0 github.com/aws/aws-sdk-go-v2/config v1.29.1 github.com/aws/aws-sdk-go-v2/service/s3 v1.73.2 - github.com/go-sql-driver/mysql v1.10.0 + github.com/block/mysql v0.0.0-20260906201522-a3178f8dca69 github.com/siddontang/loggers v1.0.3 github.com/sirupsen/logrus v1.9.3 - github.com/stretchr/testify v1.12.0 + github.com/stretchr/testify v1.12.1 golang.org/x/sync v0.22.0 ) require ( filippo.io/edwards25519 v1.2.0 // indirect - github.com/block/spirit v0.16.1-0.20260818201143-cc006c921a50 - github.com/rogpeppe/go-internal v1.11.0 // indirect + github.com/block/spirit v0.17.1-0.20260906214441-dc3d4c9f4c3b ) require ( @@ -46,21 +45,20 @@ require ( github.com/klauspost/asmfmt v1.3.2 // indirect github.com/klauspost/compress v1.18.6 // indirect github.com/klauspost/cpuid/v2 v2.2.9 // indirect - github.com/kr/pretty v0.3.0 // indirect github.com/minio/asm2plan9s v0.0.0-20200509001527-cdd76441f9d8 // indirect github.com/minio/c2goasm v0.0.0-20190812172519-36a3d3bbc4f3 // indirect github.com/pierrec/lz4/v4 v4.1.22 // indirect github.com/zeebo/xxh3 v1.0.2 // indirect + go.yaml.in/yaml/v3 v3.0.5 // indirect golang.org/x/exp v0.0.0-20240909161429-701f63a606c0 // indirect - golang.org/x/mod v0.37.0 // indirect - golang.org/x/net v0.56.0 // indirect - golang.org/x/sys v0.46.0 // indirect - golang.org/x/telemetry v0.0.0-20260625142307-59b4966ccb57 // indirect - golang.org/x/text v0.40.0 // indirect - golang.org/x/tools v0.47.0 // indirect + golang.org/x/mod v0.38.0 // indirect + golang.org/x/net v0.57.0 // indirect + golang.org/x/sys v0.47.0 // indirect + golang.org/x/telemetry v0.0.0-20260708182218-49f421fb7959 // indirect + golang.org/x/text v0.41.0 // indirect + golang.org/x/tools v0.48.0 // indirect golang.org/x/xerrors v0.0.0-20240903120638-7835f813f4da // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20241104194629-dd2ea8efbc28 // indirect google.golang.org/grpc v1.69.2 // indirect google.golang.org/protobuf v1.36.1 // indirect - gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/go.sum b/go.sum index 3630407..92fbde4 100644 --- a/go.sum +++ b/go.sum @@ -50,17 +50,16 @@ github.com/aws/aws-sdk-go-v2/service/sts v1.33.9 h1:BRVDbewN6VZcwr+FBOszDKvYeXY1 github.com/aws/aws-sdk-go-v2/service/sts v1.33.9/go.mod h1:f6vjfZER1M17Fokn0IzssOTMT2N8ZSq+7jnNF0tArvw= github.com/aws/smithy-go v1.22.1 h1:/HPHZQ0g7f4eUeK6HKglFz8uwVfZKgoI25rb/J+dnro= github.com/aws/smithy-go v1.22.1/go.mod h1:irrKGvNn1InZwb2d7fkIRNucdfwR8R+Ts3wxYa/cJHg= -github.com/block/spirit v0.16.1-0.20260818201143-cc006c921a50 h1:5ZpWwNmXYgeX4Bw2L11QssKhK902KvBwmayfvdyUFX4= -github.com/block/spirit v0.16.1-0.20260818201143-cc006c921a50/go.mod h1:tARvFXx8EfoMZv3f6gcicwNvy6/QQT1LQTWc5LrCuHA= -github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= +github.com/block/mysql v0.0.0-20260906201522-a3178f8dca69 h1:rCWVZKT5PdrdfosvMavnnfUBAdtaM3iepvADt5pMSJI= +github.com/block/mysql v0.0.0-20260906201522-a3178f8dca69/go.mod h1:KEo73lbxXs9cFlq+x3Z35UqGg3MTxAPfjDOR/ob/iik= +github.com/block/spirit v0.17.1-0.20260906214441-dc3d4c9f4c3b h1:8WteKJWigls5NasLxwqibivSa2mob4Zt4+6+wixTZXM= +github.com/block/spirit v0.17.1-0.20260906214441-dc3d4c9f4c3b/go.mod h1:x0MOoDU9krsDsdYbfEqREfrI4GLiy/7l5SGSbKsSgH4= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/go-logr/logr v1.4.2 h1:6pFjapn8bFcIbiKo3XT4j/BhANplGihG6tvd+8rYgrY= github.com/go-logr/logr v1.4.2/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= -github.com/go-sql-driver/mysql v1.10.0 h1:Q+1LV8DkHJvSYAdR83XzuhDaTykuDx0l6fkXxoWCWfw= -github.com/go-sql-driver/mysql v1.10.0/go.mod h1:M+cqaI7+xxXGG9swrdeUIoPG3Y3KCkF0pZej+SK+nWk= github.com/goccy/go-json v0.10.6 h1:p8HrPJzOakx/mn/bQtjgNjdTcN+/S6FcG2CTtQOrHVU= github.com/goccy/go-json v0.10.6/go.mod h1:oq7eo15ShAhp70Anwd5lgX2pLfOS3QCiwU/PULtXL6M= github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= @@ -81,13 +80,6 @@ github.com/klauspost/compress v1.18.6 h1:2jupLlAwFm95+YDR+NwD2MEfFO9d4z4Prjl1XXD github.com/klauspost/compress v1.18.6/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= github.com/klauspost/cpuid/v2 v2.2.9 h1:66ze0taIn2H33fBvCkXuv9BmCwDfafmiIVpKV9kKGuY= github.com/klauspost/cpuid/v2 v2.2.9/go.mod h1:rqkxqrZ1EhYM9G+hXH7YdowN5R5RGN6NK4QwQ3WMXF8= -github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= -github.com/kr/pretty v0.3.0 h1:WgNl7dwNpEZ6jJ9k1snq4pZsg7DOEN8hP9Xw0Tsjwk0= -github.com/kr/pretty v0.3.0/go.mod h1:640gp4NfQd8pI5XOwp5fnNeVWj67G7CFk/SaSQn7NBk= -github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ= -github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI= -github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= -github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= github.com/minio/asm2plan9s v0.0.0-20200509001527-cdd76441f9d8 h1:AMFGa4R4MiIpspGNG7Z948v4n35fFGB3RR3G/ry4FWs= github.com/minio/asm2plan9s v0.0.0-20200509001527-cdd76441f9d8/go.mod h1:mC1jAcsrzbxHt8iiaC+zU4b1ylILSosueou12R++wfY= github.com/minio/c2goasm v0.0.0-20190812172519-36a3d3bbc4f3 h1:+n/aFZefKZp7spd8DFdX7uMikMLXX4oubIzJF4kv/wI= @@ -95,9 +87,6 @@ github.com/minio/c2goasm v0.0.0-20190812172519-36a3d3bbc4f3/go.mod h1:RagcQ7I8Ie github.com/pierrec/lz4/v4 v4.1.22 h1:cKFw6uJDK+/gfw5BcDL0JL5aBsAFdsIT18eRtLj7VIU= github.com/pierrec/lz4/v4 v4.1.22/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= -github.com/rogpeppe/go-internal v1.6.1/go.mod h1:xXDCJY+GAPziupqXw64V24skbSoqbTEfhy4qGm1nDQc= -github.com/rogpeppe/go-internal v1.11.0 h1:cWPaGQEPrBb5/AsnsZesgZZ9yb1OQ+GOISoDNXVBh4M= -github.com/rogpeppe/go-internal v1.11.0/go.mod h1:ddIwULY96R17DhadqLgMfk9H9tvdUzkipdSkR5nkCZA= github.com/siddontang/loggers v1.0.3 h1:UzEq/9uIvze+OnZsI50FZi9q79HUaRWcD3ZuIWIOu6Y= github.com/siddontang/loggers v1.0.3/go.mod h1:22jUdLiNuifMxKFmskmzVNGx0WdprmNctzWgn0KDrCU= github.com/sirupsen/logrus v1.9.3 h1:dueUQJ1C2q9oE3F7wvmSGAaVtTmUizReu6fjN8uqzbQ= @@ -106,8 +95,8 @@ github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+ github.com/stretchr/objx v0.5.3 h1:jmXUvGomnU1o3W/V5h2VEradbpJDwGrzugQQvL0POH4= github.com/stretchr/objx v0.5.3/go.mod h1:rDQraq+vQZU7Fde9LOZLr8Tax6zZvy4kuNKF+QYS+U0= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= -github.com/stretchr/testify v1.12.0 h1:K6Mr6jO9JICuend/5xzTM03ydSV3vdNRYAdPSukj8uI= -github.com/stretchr/testify v1.12.0/go.mod h1:bOYBZb5qJ00vPzWfIqBUZPaxK8jWiXc6d3ErP4Ca9Gw= +github.com/stretchr/testify v1.12.1 h1:EuwCh5fleGS7H32xRwO3wRGT7DxrDhLAT6FF8MpWDWE= +github.com/stretchr/testify v1.12.1/go.mod h1:MDEgiDPPsNp5cuIrHPPCyornHKgEVbtFUmoNlxoYthg= github.com/xyproto/randomstring v1.0.5 h1:YtlWPoRdgMu3NZtP45drfy1GKoojuR7hmRcnhZqKjWU= github.com/xyproto/randomstring v1.0.5/go.mod h1:rgmS5DeNXLivK7YprL0pY+lTuhNQW3iGxZ18UQApw/E= github.com/zeebo/assert v1.3.0 h1:g7C04CbJuIDKNPFHmsk4hwZDO5O+kntRxzaUoNXj+IQ= @@ -126,23 +115,25 @@ go.opentelemetry.io/otel/trace v1.31.0 h1:ffjsj1aRouKewfr85U2aGagJ46+MvodynlQ1HY go.opentelemetry.io/otel/trace v1.31.0/go.mod h1:TXZkRk7SM2ZQLtR6eoAWQFIHPvzQ06FJAsO1tJg480A= go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= +go.yaml.in/yaml/v3 v3.0.5 h1:N6y/pJk8buWs9NY5ERU2HSMfm+IuD/OtfdAnq6kESPw= +go.yaml.in/yaml/v3 v3.0.5/go.mod h1:HVTZu1O7/Vkt2N+BFy8Zza+lnLsABggaTM2ZpNIGuKg= golang.org/x/exp v0.0.0-20240909161429-701f63a606c0 h1:e66Fs6Z+fZTbFBAxKfP3PALWBtpfqks2bwGcexMxgtk= golang.org/x/exp v0.0.0-20240909161429-701f63a606c0/go.mod h1:2TbTHSBQa924w8M6Xs1QcRcFwyucIwBGpK1p2f1YFFY= -golang.org/x/mod v0.37.0 h1:vF1DjpVEshcIqoEaauuHebaLk1O1forxjxBaVn884JQ= -golang.org/x/mod v0.37.0/go.mod h1:m8S8VeM9r4dzDwjrKO0a1sZP3YjeMamRRlD+fmR2Q/0= -golang.org/x/net v0.56.0 h1:Rw8j/hFzGvJUZwNBXnAtf5sVDVt+65SK2C7IxCxZt5o= -golang.org/x/net v0.56.0/go.mod h1:D3Ku6r+V6JROoZK144D2XfMHFcMq/0zSfLelVTCFKec= +golang.org/x/mod v0.38.0 h1:MECBjubtXD7yj4HrhIUcywNaGeNVUdfVnxmPajOk4yk= +golang.org/x/mod v0.38.0/go.mod h1:V6Xz0pq8TQ3dGqVQ1FVHuelZpAL0uNhSkk9ogYP3c40= +golang.org/x/net v0.57.0 h1:K5+3DljvIuDG9/Jv9rvyMywYNFCQ9RSUY6OOTTkT+tE= +golang.org/x/net v0.57.0/go.mod h1:KpXc8iv+r3XplLAG/f7Jsf9RPszJzdR0f58q9vGOuEU= golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek= golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= golang.org/x/sys v0.0.0-20220715151400-c0bba94af5f8/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.46.0 h1:noSf2Fq6F8DBgS+LysIkx7rIExoNHJsxOAtPp4rthXw= -golang.org/x/sys v0.46.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= -golang.org/x/telemetry v0.0.0-20260625142307-59b4966ccb57 h1:nwGZBCt+FnXUrGsj5vjzAsEmkcaFvd82BbOjECiFYZc= -golang.org/x/telemetry v0.0.0-20260625142307-59b4966ccb57/go.mod h1:3AWMyWHS+caVoiEXpiq6+tzKA40J4vQT3MYr80ZtQpc= -golang.org/x/text v0.40.0 h1:Ub2Z6/xjgF1WrYQz2nuITOEegKFtiIy+rieRJ5lHZKs= -golang.org/x/text v0.40.0/go.mod h1:hpnzDAfGV753zIKo+wk3u1bVKCGPbrnF7+7LBF/UHVY= -golang.org/x/tools v0.47.0 h1:7Kn5x/d1svx/PzryTsqeoZN4TZwqeH5pGWjefhLi/1Q= -golang.org/x/tools v0.47.0/go.mod h1:dFHnyTvFWY212G+h7ZY4Vsp/K3U4/7W9TyVaAul8uCA= +golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs= +golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/telemetry v0.0.0-20260708182218-49f421fb7959 h1:RJhm5l6Fo4rmEIcndxDllNhhf/fAx8qIm4t6A7vpm2A= +golang.org/x/telemetry v0.0.0-20260708182218-49f421fb7959/go.mod h1:LV7u5Oco+Z/g6XI7PqN+EUUUGGkEcmB1uj2ceI0fOVg= +golang.org/x/text v0.41.0 h1:vz/seA0lnX87Othu2f/0L24RcgrXD9/YFTSuGjj3rH8= +golang.org/x/text v0.41.0/go.mod h1:jvf1O8ajNzZqhSrQBPbutR/EB83Cc0CFrezNQIwbb5M= +golang.org/x/tools v0.48.0 h1:3+hClM1aLL5mjMKm5ovokw9epgRXPuu2tILgismM6RE= +golang.org/x/tools v0.48.0/go.mod h1:08xX0orndb/F7jJxGDicx061tyd5pcMto75YMAXr6lk= golang.org/x/xerrors v0.0.0-20240903120638-7835f813f4da h1:noIWHXmPHxILtqtCOPIhSt0ABwskkZKjD3bXGnZGpNY= golang.org/x/xerrors v0.0.0-20240903120638-7835f813f4da/go.mod h1:NDW/Ps6MPRej6fsCIbMTohpP40sJ/P/vI1MoTEGwX90= gonum.org/v1/gonum v0.15.1 h1:FNy7N6OUZVUaWG9pTiD+jlhdQ3lMP+/LcTpJ6+a8sQ0= @@ -154,10 +145,4 @@ google.golang.org/grpc v1.69.2/go.mod h1:vyjdE6jLBI76dgpDojsFGNaHlxdjXN9ghpnd2o7 google.golang.org/protobuf v1.36.1 h1:yBPeRvTftaleIgM3PZ/WBIZ7XM/eEYAaEyCwvyjq/gk= google.golang.org/protobuf v1.36.1/go.mod h1:9fA7Ob0pmnwhb644+1+CVWFRbNajQ6iRojtC/QF5bRE= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= -gopkg.in/check.v1 v1.0.0-20180628173108-788fd7840127/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= -gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= -gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= -gopkg.in/errgo.v2 v2.1.0/go.mod h1:hNsd1EY+bozCKY1Ytp96fpM3vjJbqLJn88ws8XvfDNI= gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= -gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= -gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/pkg/archive/buffer_stager_test.go b/pkg/archive/buffer_stager_test.go index ca94b24..780af1c 100644 --- a/pkg/archive/buffer_stager_test.go +++ b/pkg/archive/buffer_stager_test.go @@ -6,11 +6,11 @@ import ( "testing" "time" + "github.com/block/mysql" "github.com/block/polt/pkg/parquet" "github.com/block/polt/pkg/stage" "github.com/block/polt/pkg/test" "github.com/block/spirit/pkg/table" - "github.com/go-sql-driver/mysql" "github.com/sirupsen/logrus" "github.com/stretchr/testify/require" ) @@ -20,7 +20,7 @@ func TestNewBufferStager(t *testing.T) { cfg, err := mysql.ParseDSN(test.DSN()) require.NoError(t, err) - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) tbl := `CREATE TABLE _t1_ia_stage_buffer_runid ( diff --git a/pkg/archive/file_test.go b/pkg/archive/file_test.go index e8dd61b..e470aeb 100644 --- a/pkg/archive/file_test.go +++ b/pkg/archive/file_test.go @@ -10,10 +10,10 @@ import ( "github.com/apache/arrow-go/v18/arrow/memory" "github.com/apache/arrow-go/v18/parquet" "github.com/apache/arrow-go/v18/parquet/file" + "github.com/block/mysql" "github.com/block/polt/pkg/test" "github.com/block/polt/pkg/upload" "github.com/block/spirit/pkg/table" - "github.com/go-sql-driver/mysql" "github.com/sirupsen/logrus" "github.com/stretchr/testify/require" ) @@ -42,7 +42,7 @@ func TestNewFileArchiver(t *testing.T) { cfg, err := mysql.ParseDSN(test.DSN()) require.NoError(t, err) - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) stgTableCnt := test.GetCount(t, db, "_t1_ia_stage_file_runid", "1=1") @@ -114,7 +114,7 @@ func TestNewFileArchiver_BinaryKey(t *testing.T) { cfg, err := mysql.ParseDSN(test.DSN()) require.NoError(t, err) - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) stgTableCnt := test.GetCount(t, db, "_t2_ia_stage_file_runid", "1=1") diff --git a/pkg/archive/table_test.go b/pkg/archive/table_test.go index 015b576..96d560f 100644 --- a/pkg/archive/table_test.go +++ b/pkg/archive/table_test.go @@ -5,8 +5,8 @@ import ( "database/sql" "testing" + "github.com/block/mysql" "github.com/block/spirit/pkg/table" - "github.com/go-sql-driver/mysql" "github.com/stretchr/testify/require" "github.com/block/polt/pkg/test" @@ -31,7 +31,7 @@ func TestTableArchiver_Move(t *testing.T) { cfg, err := mysql.ParseDSN(test.DSN()) require.NoError(t, err) - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) stagedCount := test.GetCount(t, db, "_t1_ia_stage_runid", "1=1") diff --git a/pkg/audit/db_test.go b/pkg/audit/db_test.go index bf6d40a..f108f4b 100644 --- a/pkg/audit/db_test.go +++ b/pkg/audit/db_test.go @@ -12,7 +12,7 @@ import ( func TestCreateAuditDB(t *testing.T) { // Create a new audit database // Check that the audit database is created correctly - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) defer func() { if closeErr := db.Close(); closeErr != nil { @@ -30,7 +30,7 @@ func TestCreateAuditDB(t *testing.T) { func TestCreateCheckpointTbl(t *testing.T) { // Create a new checkpoint table // Check that the checkpoint table is created correctly - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) defer func() { if closeErr := db.Close(); closeErr != nil { @@ -50,7 +50,7 @@ func TestCreateCheckpointTbl(t *testing.T) { func TestCreateRunsTbl(t *testing.T) { // Create a new runs table // Check that the runs table is created correctly - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) defer func() { if closeErr := db.Close(); closeErr != nil { diff --git a/pkg/boot/archive_booter_test.go b/pkg/boot/archive_booter_test.go index 390b182..b4346a6 100644 --- a/pkg/boot/archive_booter_test.go +++ b/pkg/boot/archive_booter_test.go @@ -5,9 +5,9 @@ import ( "database/sql" "testing" + "github.com/block/mysql" "github.com/block/polt/pkg/test" "github.com/block/spirit/pkg/table" - "github.com/go-sql-driver/mysql" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) @@ -24,7 +24,7 @@ func TestArchiveBooter_Setup(t *testing.T) { cfg, err := mysql.ParseDSN(test.DSN()) require.NoError(t, err) - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) srcTbl := table.NewTableInfo(db, cfg.DBName, "_t1_sb_stage_abid") err = srcTbl.SetInfo(context.Background()) @@ -63,7 +63,7 @@ func TestArchiveBooter_PostSetupChecks(t *testing.T) { cfg, err := mysql.ParseDSN(test.DSN()) require.NoError(t, err) - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) srcTbl := table.NewTableInfo(db, cfg.DBName, "t1_ab") err = srcTbl.SetInfo(context.Background()) diff --git a/pkg/boot/booter_test.go b/pkg/boot/booter_test.go index 316cc9b..5b848a1 100644 --- a/pkg/boot/booter_test.go +++ b/pkg/boot/booter_test.go @@ -3,8 +3,8 @@ package boot import ( "testing" + "github.com/block/mysql" "github.com/block/polt/pkg/test" - "github.com/go-sql-driver/mysql" "github.com/stretchr/testify/require" ) diff --git a/pkg/boot/stage_booter_test.go b/pkg/boot/stage_booter_test.go index 9454ca8..bbabee8 100644 --- a/pkg/boot/stage_booter_test.go +++ b/pkg/boot/stage_booter_test.go @@ -7,10 +7,10 @@ import ( "net/url" "testing" + "github.com/block/mysql" "github.com/block/polt/pkg/query" "github.com/block/polt/pkg/test" "github.com/block/spirit/pkg/table" - "github.com/go-sql-driver/mysql" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) @@ -28,7 +28,7 @@ func TestStageBooter_Setup(t *testing.T) { cfg, err := mysql.ParseDSN(test.DSN()) require.NoError(t, err) - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) srcTbl := table.NewTableInfo(db, cfg.DBName, "t1_sb") err = srcTbl.SetInfo(context.Background()) @@ -59,7 +59,7 @@ func TestStageBooter_PreflightChecks(t *testing.T) { cfg, err := mysql.ParseDSN(test.DSN()) require.NoError(t, err) - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) srcTbl := table.NewTableInfo(db, cfg.DBName, "t1_sb") err = srcTbl.SetInfo(context.Background()) @@ -92,9 +92,9 @@ func TestStageBooter_PreflightChecks_ReplicaHealth(t *testing.T) { cfg, err := mysql.ParseDSN(test.DSN()) require.NoError(t, err) - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) - replicadb, err := sql.Open("mysql", test.ReplicaDSN()) + replicadb, err := sql.Open("block-mysql", test.ReplicaDSN()) require.NoError(t, err) srcTbl := table.NewTableInfo(db, cfg.DBName, "t1_sb") err = srcTbl.SetInfo(context.Background()) @@ -105,7 +105,7 @@ func TestStageBooter_PreflightChecks_ReplicaHealth(t *testing.T) { // use a completely invalid DSN. // golang sql.Open lazy loads, so this is possible. - replicadb, err = sql.Open("mysql", "msandbox:msandbox@tcp(127.0.0.1:22)/test") + replicadb, err = sql.Open("block-mysql", "msandbox:msandbox@tcp(127.0.0.1:22)/test") require.NoError(t, err) s = NewStageBooter(&StageBooterConfig{AuditDB: "polt", Query: "SELECT * from t1_sb WHERE name = 'harry'", RunID: "runid", DB: db, SrcTbl: srcTbl, Replica: replicadb}) @@ -133,7 +133,7 @@ func TestStageBooter_PreflightChecks_PreEvalQuery(t *testing.T) { // For the purpose of this test SET TIMESTAMP = 1696911031 and SET TIME_ZONE='UTC', // which returns the fixed value for NOW() as '2023-10-10 04:10:31' timeZone := fmt.Sprintf("%s=%s", "time_zone", url.QueryEscape(`"+00:00"`)) - db, err := sql.Open("mysql", test.DSN()+"?timestamp=1696911031&"+timeZone) + db, err := sql.Open("block-mysql", test.DSN()+"?timestamp=1696911031&"+timeZone) require.NoError(t, err) srcTbl := table.NewTableInfo(db, cfg.DBName, "t1_sb") diff --git a/pkg/parquet/arrow_schema_test.go b/pkg/parquet/arrow_schema_test.go index 712a49f..304be6e 100644 --- a/pkg/parquet/arrow_schema_test.go +++ b/pkg/parquet/arrow_schema_test.go @@ -21,7 +21,7 @@ func TestMysqlToArrowSchema(t *testing.T) { PRIMARY KEY (id) )` test.RunSQL(t, tbl) - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) aSchema, err := MysqlToArrowSchema(db, "schema_test") require.NoError(t, err) diff --git a/pkg/parquet/write_buffer_test.go b/pkg/parquet/write_buffer_test.go index 6034130..ebd26e5 100644 --- a/pkg/parquet/write_buffer_test.go +++ b/pkg/parquet/write_buffer_test.go @@ -93,7 +93,7 @@ func TestNewParquetManagerFlush(t *testing.T) { test.RunSQL(t, "INSERT INTO t_flush (id) VALUES (2)") test.RunSQL(t, "INSERT INTO t_flush (id) VALUES (3)") - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) ti := table.NewTableInfo(db, "test", "t_flush") err = ti.SetInfo(context.Background()) diff --git a/pkg/query/explain_test.go b/pkg/query/explain_test.go index 08a777e..ab24dc7 100644 --- a/pkg/query/explain_test.go +++ b/pkg/query/explain_test.go @@ -18,7 +18,7 @@ func Test_Explain(t *testing.T) { )` test.RunSQL(t, tbl) - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) query := "SELECT * FROM t1_q WHERE name = 'harry'" diff --git a/pkg/query/expression_folder_test.go b/pkg/query/expression_folder_test.go index a95ee5e..f1637c8 100644 --- a/pkg/query/expression_folder_test.go +++ b/pkg/query/expression_folder_test.go @@ -7,9 +7,9 @@ import ( "net/url" "testing" + "github.com/block/mysql" "github.com/block/polt/pkg/test" "github.com/block/spirit/pkg/table" - "github.com/go-sql-driver/mysql" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) @@ -30,7 +30,7 @@ func Test_PreEvaluate(t *testing.T) { // For the purpose of this test SET TIMESTAMP = 1696911031 and SET TIME_ZONE='UTC', // which returns the fixed value for NOW() as '2023-10-10 04:10:31' timeZone := fmt.Sprintf("%s=%s", "time_zone", url.QueryEscape(`"+00:00"`)) - db, err := sql.Open("mysql", test.DSN()+"?timestamp=1696911031&"+timeZone) + db, err := sql.Open("block-mysql", test.DSN()+"?timestamp=1696911031&"+timeZone) require.NoError(t, err) srcTbl := table.NewTableInfo(db, cfg.DBName, "t1_q") diff --git a/pkg/query/validator_test.go b/pkg/query/validator_test.go index 64de351..4b3a542 100644 --- a/pkg/query/validator_test.go +++ b/pkg/query/validator_test.go @@ -5,9 +5,9 @@ import ( "database/sql" "testing" + "github.com/block/mysql" "github.com/block/polt/pkg/test" "github.com/block/spirit/pkg/table" - "github.com/go-sql-driver/mysql" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) @@ -23,7 +23,7 @@ func Test_Validate(t *testing.T) { cfg, err := mysql.ParseDSN(test.DSN()) require.NoError(t, err) - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) srcTbl := table.NewTableInfo(db, cfg.DBName, "t1_q") err = srcTbl.SetInfo(context.Background()) diff --git a/pkg/runner/archive_runner_test.go b/pkg/runner/archive_runner_test.go index 365c5a8..5e7ad07 100644 --- a/pkg/runner/archive_runner_test.go +++ b/pkg/runner/archive_runner_test.go @@ -11,9 +11,9 @@ import ( "github.com/apache/arrow-go/v18/arrow/memory" "github.com/apache/arrow-go/v18/parquet" "github.com/apache/arrow-go/v18/parquet/file" + "github.com/block/mysql" "github.com/block/polt/pkg/destinations" "github.com/block/polt/pkg/test" - "github.com/go-sql-driver/mysql" "github.com/siddontang/loggers" "github.com/sirupsen/logrus" "github.com/stretchr/testify/assert" @@ -300,7 +300,7 @@ func TestArchiveMultipleParquetFiles(t *testing.T) { name2 varchar(255) NOT NULL )` - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) test.RunSQL(t, tbl) @@ -383,7 +383,7 @@ func TestArchiveMultipleParquetFiles_ResumeFromCheckpt(t *testing.T) { name2 varchar(255) NOT NULL )` - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) test.RunSQL(t, tbl) @@ -551,7 +551,7 @@ func TestArchiveRunnerWithGeneratedColumns(t *testing.T) { require.NoError(t, a.Run(context.Background())) // Verify that the archive table was created - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) defer func() { if closeErr := db.Close(); closeErr != nil { diff --git a/pkg/runner/runner_test.go b/pkg/runner/runner_test.go index b8d4a01..6addf11 100644 --- a/pkg/runner/runner_test.go +++ b/pkg/runner/runner_test.go @@ -15,7 +15,7 @@ func TestRunEntry(t *testing.T) { // Create a new run entry // Check that the run entry is created correctly test.RunSQL(t, "DROP TABLE IF EXISTS polt.runs") - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) defer func() { if closeErr := db.Close(); closeErr != nil { diff --git a/pkg/runner/stage_runner_test.go b/pkg/runner/stage_runner_test.go index 9b4ffbd..1dfc17b 100644 --- a/pkg/runner/stage_runner_test.go +++ b/pkg/runner/stage_runner_test.go @@ -6,8 +6,8 @@ import ( "testing" "time" + "github.com/block/mysql" "github.com/block/polt/pkg/test" - "github.com/go-sql-driver/mysql" "github.com/siddontang/loggers" "github.com/sirupsen/logrus" "github.com/stretchr/testify/assert" @@ -316,7 +316,7 @@ func TestResumeFromCheckpoint(t *testing.T) { INDEX(idxed_column) )` - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) test.RunSQL(t, tbl) @@ -454,7 +454,7 @@ func TestStageRunnerWithGeneratedColumns(t *testing.T) { require.Equal(t, "_t1_sr_gen_stage_runid", stgTbl) // Verify that the staging table was created and data was copied - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) defer func() { if closeErr := db.Close(); closeErr != nil { @@ -514,7 +514,7 @@ func TestResumeFromDBFailure(t *testing.T) { INDEX(idxed_column) )` - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) test.RunSQL(t, tbl) diff --git a/pkg/stage/stager_test.go b/pkg/stage/stager_test.go index a18c5a8..1d90a1d 100644 --- a/pkg/stage/stager_test.go +++ b/pkg/stage/stager_test.go @@ -6,10 +6,10 @@ import ( "testing" "time" + "github.com/block/mysql" "github.com/block/polt/pkg/test" "github.com/block/spirit/pkg/table" "github.com/block/spirit/pkg/throttler" - "github.com/go-sql-driver/mysql" "github.com/sirupsen/logrus" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -50,7 +50,7 @@ func TestStager_DryRun_Stage(t *testing.T) { setup(t) cfg, err := mysql.ParseDSN(test.DSN()) require.NoError(t, err) - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) srcTbl := table.NewTableInfo(db, cfg.DBName, "t1_st") @@ -95,7 +95,7 @@ func TestStager_Stage(t *testing.T) { cfg, err := mysql.ParseDSN(test.DSN()) require.NoError(t, err) - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) srcTbl := table.NewTableInfo(db, cfg.DBName, "t1_st") @@ -164,7 +164,7 @@ func TestETA(t *testing.T) { cfg, err := mysql.ParseDSN(test.DSN()) require.NoError(t, err) - db, err := sql.Open("mysql", test.DSN()) + db, err := sql.Open("block-mysql", test.DSN()) require.NoError(t, err) srcTbl := table.NewTableInfo(db, cfg.DBName, "t1_st") diff --git a/pkg/test/test.go b/pkg/test/test.go index 9895e4c..c47209e 100644 --- a/pkg/test/test.go +++ b/pkg/test/test.go @@ -9,9 +9,9 @@ import ( "path/filepath" "testing" + "github.com/block/mysql" "github.com/block/spirit/pkg/dbconn" "github.com/block/spirit/pkg/table" - "github.com/go-sql-driver/mysql" "github.com/stretchr/testify/require" ) @@ -29,7 +29,7 @@ func ReplicaDSN() string { func RunSQL(t *testing.T, stmt string) { t.Helper() - db, err := sql.Open("mysql", DSN()) + db, err := sql.Open("block-mysql", DSN()) require.NoError(t, err) defer func() { if closeErr := db.Close(); closeErr != nil {