1、Flink 2.1:SQL Advances for Data+AIBridging real-time data processing and AI in Flink SQLLincoln Lee Alibaba CloudStaff Engineer,Apache Flink PMCLLM UDFIntegrate AI Analytics into Real-Time PipelineHuman Reviewtable product(id BIGINT,title VARCHAR,country INT category VARCHAR)table risk_output(id BIG
2、INT,title VARCHAR,forbidden_words VARCHAR hit_keywords VARCHAR)Product title:grape juice(no alcohol)LLMUnderstand no alcoholMeans doesnt contain alcoholno riskclass OpenAIUDF extends ScalarFunction void open()/initOpenAIClient.String eval(String title)/1.prepare ChatCompletionCreateParams/2.return c
3、lient.chat().completions().create.The Hidden Costs of Custom AI UDFLLM UDFOpenAItry different models 1 seconds3 secondssynchronize requestlow throughput async request higher throughputScalarFunctionAsyncScalarFunctionrewrite coderewrite codeApache Flink SQL Native AI FunctionCreate ModelCREATE MODEL
4、 my_modelINPUT(input STRING)OUTPUT(content STRING)WITH(provider=openai/modelScope/,model=gpt-4o/deepseek/qwen/,endpoint=https:/,systemPrompt=)Chat/CompletionSELECT *,content as ai_result FROM ML_PREDICT(TABLE user_feedbacks,MODELmy_model,DESCRIPTOR(feedback),MAPasync,true)SELECT *,content as feature
5、sFROM ML_PREDICT(TABLE input,MODELmy_embedding_model,DESCRIPTOR(question),MAPasync,true)Embeddingmodel managementSQL AI function JSON is important for AIJSON Everywhere:From Big Data to AISearch documentEvent logsJSON_QUERYJSON_VALUEJSON_STRINGJSON_OBJECTJSON_EXISTSJSON_ARRAYSemi-structured JSON Dat
6、aBuilt-in JSON FunctionsProcess JSON StringRAG documentLLM prompts/outputAPI payload JSONJSONparsing”events The Hidden Cost of JSON ParsinguserId:u123,events:type:click,ts:1717230000,type:purchase,ts:1717230033,amount:59.9,metadata:browser:Chrome,device:MobileJSON DocumentRequires parsing every time