Hi :wave: I'm trying to find the best way to deal ...
# troubleshooting
s
Hi 👋 I'm trying to find the best way to deal with deleted records. I have Mysql + debezium + kafka + pinot setup. When a record is deleted from the source db, the
after
part of the payload on kafka will be null as expected. I am extracting the fields from
$.after
of the payload (attached my table def and schema def). But when the record is deleted, I want to mark the deleted record with a
deleted_at
column and extract the fields from the
$.before
part of the payload. I haven't used Groovy before and not sure if I'm doing it right, as a test I'm trying to set
my_id
field to payload.after.id if after is not null and otherwise set it to payload.before.id:
Copy code
{
        "columnName": "my_id",
        "transformFunction": "Groovy({JSONPATHLONG(payload, '$.after.id', '0') == 0 ? JSONPATHLONG(payload, '$.before.id', '0') : JSONPATHLONG(payload, '$.after.id', '0')}, payload)"
      },
but it fails creating the table with a 400 invalid table config error. Any help/hints would be appreciated or how others deal with deleted records coming from source
def-schema.json,def-table.json
m
@Rong R ^^
m
I'm not sure if you can use the JSON functions inside a Groovy script - you might want to instead use Groovy syntax to do the coercion of the
id
field to a long
r
Yes. Mark is correct, JSONPATHLONG is a built-in ScalarFunction in Pinot that cannot be understood by Groovy shell.
s
thanks all. I ended up replacing the transformation in the table definition with
Copy code
{
        "columnName": "my_id",
        "transformFunction": "Groovy('{"returnType":"LONG","isSingleValue":true}'
                                     'def jsonSlurper = new JsonSlurper();
                                     def object = jsonSlurper.parseText(payload);
                                     def result = 0;
                                     {object.after == null ? {result=Long.valueOf(object.before.id)} : {result=Long.valueOf(object.after.id}};
                                     return result',payload)"
      },
but I am getting
Copy code
shaded.com.fasterxml.jackson.core.JsonParseException: Unexpected character ('r' (code 114)): was expecting comma to separate Object entries
 at [Source: (FileInputStream); line: 44, column: 42]
	at shaded.com.fasterxml.jackson.core.JsonParser._constructError(JsonParser.java:1840)
	at shaded.com.fasterxml.jackson.core.base.ParserMinimalBase._reportError(ParserMinimalBase.java:712)
	at shaded.com.fasterxml.jackson.core.base.ParserMinimalBase._reportUnexpectedChar(ParserMinimalBase.java:637)
	at shaded.com.fasterxml.jackson.core.json.UTF8StreamJsonParser.nextFieldName(UTF8StreamJsonParser.java:1010)
	at shaded.com.fasterxml.jackson.databind.deser.std.BaseNodeDeserializer.deserializeObject(JsonNodeDeserializer.java:250)
	at shaded.com.fasterxml.jackson.databind.deser.std.BaseNodeDeserializer.deserializeArray(JsonNodeDeserializer.java:437)
	at shaded.com.fasterxml.jackson.databind.deser.std.BaseNodeDeserializer.deserializeObject(JsonNodeDeserializer.java:261)
	at shaded.com.fasterxml.jackson.databind.deser.std.BaseNodeDeserializer.deserializeObject(JsonNodeDeserializer.java:258)
	at shaded.com.fasterxml.jackson.databind.deser.std.JsonNodeDeserializer.deserialize(JsonNodeDeserializer.java:68)
	at shaded.com.fasterxml.jackson.databind.deser.std.JsonNodeDeserializer.deserialize(JsonNodeDeserializer.java:15)
	at shaded.com.fasterxml.jackson.databind.ObjectReader._bindAsTree(ObjectReader.java:1770)
	at shaded.com.fasterxml.jackson.databind.ObjectReader._bindAndCloseAsTree(ObjectReader.java:1735)
	at shaded.com.fasterxml.jackson.databind.ObjectReader.readTree(ObjectReader.java:1395)
	at org.apache.pinot.spi.utils.JsonUtils.fileToJsonNode(JsonUtils.java:104)
	at org.apache.pinot.tools.admin.command.AddTableCommand.execute(AddTableCommand.java:205)
	at org.apache.pinot.tools.Command.call(Command.java:33)
	at org.apache.pinot.tools.Command.call(Command.java:29)
	at picocli.CommandLine.executeUserObject(CommandLine.java:1953)
	at picocli.CommandLine.access$1300(CommandLine.java:145)
	at picocli.CommandLine$RunLast.executeUserObjectOfLastSubcommandWithSameParent(CommandLine.java:2352)
	at picocli.CommandLine$RunLast.handle(CommandLine.java:2346)
	at picocli.CommandLine$RunLast.handle(CommandLine.java:2311)
	at picocli.CommandLine$AbstractParseResultHandler.execute(CommandLine.java:2179)
	at picocli.CommandLine.execute(CommandLine.java:2078)
	at org.apache.pinot.tools.admin.PinotAdministrator.execute(PinotAdministrator.java:161)
	at org.apache.pinot.tools.admin.PinotAdministrator.main(PinotAdministrator.java:192)
when creating the table
r
this looks like you have some sort of JSON encoding error in. your addTable command payload. can you share the full json file (or line 44: column 42 🙂 ) thanks
s
hmm interesting, I've been using the same command (with different table config files, it works fine):
Copy code
docker exec themis-pinot-1 bash -c "/opt/pinot/bin/pinot-admin.sh AddTable -tableConfigFile /opt/pinot/def-table.json -schemaFile /opt/pinot/def-schema.json -exec"
r
yeah that's what I meant - can you share the json table config file that prompts that error. there must've been something wrong in line 44 of that file
s
ah sorry, yes,
this is the only section I've been modifying ( I think):
Copy code
{
    "columnName": "my_id",
        "transformFunction": "Groovy('{"returnType":"LONG","isSingleValue":true}'
                                     'def jsonSlurper = new JsonSlurper();
                                     def object = jsonSlurper.parseText(payload);
                     def result = 0;
                                     {object.after == null ? {result=Long.valueOf(object.before.id)} : {result=Long.valueOf(object.after.id}};
                                     return result',payload)"
      },
r
oh you had to escape the double quote
e.g.
Copy code
{
    "columnName": "my_id",
        "transformFunction": "Groovy('{\"returnType\":\"LONG\", ...
      },
s
thanks for the help and sorry for the late response, I'm getting some other errors that seems to be related having new lines in the groovy scripts. Going to fix that tomorrow and see if it works.
Hi again, I'm getting
Copy code
{
  "code": 400,
  "error": "Invalid transform function 'Groovy('{\"returnType\":\"LONG\",\"isSingleValue\":true}','def jsonSlurper = new JsonSlurper();def object = jsonSlurper.parseText(payload);def result = 0;{object.after == null ? {result=Long.valueOf(object.before.id)} : {result=Long.valueOf(object.after.id)}};return result',payload)' for column 'my_id'"
}
does my groovy look right? here is my new table def file
m
@Sahar can you try to put this in place of the Groovy transform function:
Copy code
{
  "columnName": "my_id",
  "transformFunction": "Groovy({def jsonSlurper = new groovy.json.JsonSlurper(); def object = jsonSlurper.parseText(payload); def result = object.after == null ? Long.valueOf(object.before.id) : Long.valueOf(object.after.id); return result}, payload)"
{
I tried that out using a Groovy REPL here - https://groovyconsole.appspot.com/ and it seems to work:
Copy code
def payload = '{"before": { "id": 2}, "after": { "id": 3}}'

def jsonSlurper = new groovy.json.JsonSlurper(); def object = jsonSlurper.parseText(payload); def result = object.after == null ? Long.valueOf(object.before.id) : Long.valueOf(object.after.id); return result
s
yes, it works, the table gets created successfully. However, it is failing in the parseText part and therefore no records making it from kafka to the table (table is empty)
Copy code
"fieldToValueMap" : {

    "payload" : {

      "op" : "u",

      "before" : {

        "open_date" : 19019,

        "description" : "test",

        "created_at" : 1643291274000,

        "billable" : 1,

        "client_id" : 347357,

        "number" : 53,

        "account_id" : 347321,

        "updated_at" : 1643291274000,

        "user_id" : 347321,

        "group_id" : 347321,

        "display_number" : "00053-Dickens",

        "id" : 347425,

        "status" : 1

      },

      "after" : {

        "open_date" : 19019,

        "description" : "test",

        "created_at" : 1643291274000,

        "billable" : 1,

        "client_id" : 347357,

        "number" : 53,

        "account_id" : 347321,

        "updated_at" : 1643291274000,

        "user_id" : 347321,

        "group_id" : 347323,

        "display_number" : "00053-Dickens",

        "id" : 347425,

        "status" : 1

      },

      "source" : {

        "thread" : 109182,

        "server_id" : 1,

        "version" : "1.0.0.Final",

        "file" : "docker-1-bin-log.000030",

        "connector" : "mysql",

        "pos" : 105685,

        "name" : "debezium_dev",

        "gtid" : "614c4ede-5f4e-11ec-a055-0242c0a89009:23948",

        "row" : 0,

        "ts_ms" : 1643291274000,

        "snapshot" : "false",

        "db" : "themis_development_1",

        "table" : "matters"

      },

      "ts_ms" : 1643291274311

    },

    "my_id" : null,

    "id" : null,

    "full_payload" : null,

    "ts_ms" : null

  },

  "nullValueFields" : [ ]

}

groovy.lang.MissingMethodException: No signature of method: groovy.json.JsonSlurper.parseText() is applicable for argument types: (java.util.HashMap) values: [[op:u, before:[open_date:19019, description:test, created_at:1643291274000, ...], ...]]

Possible solutions: parseText(java.lang.String), parse([B), parse([C), parse(java.io.File), parse(java.io.InputStream), parse(java.io.Reader)

at org.codehaus.groovy.runtime.ScriptBytecodeAdapter.unwrap(ScriptBytecodeAdapter.java:71) ~[pinot-all-0.10.0-SNAPSHOT-jar-with-dependencies.jar:0.10.0-SNAPSHOT-7ec47c420be9c6aee6c8e95644266fe9b7fe7a2b]

at org.codehaus.groovy.runtime.callsite.PojoMetaClassSite.call(PojoMetaClassSite.java:48) ~[pinot-all-0.10.0-SNAPSHOT-jar-with-dependencies.jar:0.10.0-SNAPSHOT-7ec47c420be9c6aee6c8e95644266fe9b7fe7a2b]

at org.codehaus.groovy.runtime.callsite.AbstractCallSite.call(AbstractCallSite.java:128) ~[pinot-all-0.10.0-SNAPSHOT-jar-with-dependencies.jar:0.10.0-SNAPSHOT-7ec47c420be9c6aee6c8e95644266fe9b7fe7a2b]

at Script1.run(Script1.groovy:1) ~[?:?]
m
oh it doesn't like the type of the payload value
so maybe need to convert the map to string first?
try this:
Copy code
{
  "columnName": "my_id",
  "transformFunction": "Groovy({def jsonSlurper = new groovy.json.JsonSlurper(); def object = jsonSlurper.parseText(new groovy.json.JsonBuilder(payload).toPrettyString()); def result = object.after == null ? Long.valueOf(object.before.id) : Long.valueOf(object.after.id); return result​}, payload)"
}
s
amazing, thanks so much for your help, it is working now and I learnt a bunch during the process. I need to learn when curly braces are needed in the groovy function and when I can just write the code with no return type and curly braces
👍 1