1.10版本 无法使用WATERMARK ,报空指针。
Ninguém assumiu esta issue ainda.
Avaliação
- Dificuldade
- 4/5
- Tempo estimado
- 3-5 dias
- Facilidade para iniciantes
- 25/100
- Tipo de issue
- Bug
- Clareza
- Razoavelmente clara
- Status de atividade
- Estagnada
- Domínio
- stream-processing
Direção de pesquisa
Comece pela instrução CREATE TABLE da issue e pelo erro CustomerWaterMarkerForLong do log do Flink 1.10. Compare o campo bigint derivado de UNIX_TIMESTAMP e a sintaxe de WATERMARK com o comportamento compatível com a versão 1.10. A tarefa estará concluída quando o schema relatado tiver um exemplo completo e funcional de WATERMARK e não produzir mais o erro de ponteiro nulo.
Escrita pelo modelo de indexação a partir do texto da issue.
Descrição
日志样例
{"time":"2020-06-17 23:00:01.211","pushid":"pushback_send-BC110-12392212311","app":"mm"}
由于时间序列是bigint类型,用UNIX_TIMESTAMP进行转换
Flink运行日志报错:
ERROR com.dtstack.flink.sql.watermarker.CustomerWaterMarkerForLong -
java.lang.NullPointerException
建表语句:
CREATE TABLE MyTable (
time varchar ,
pushid varchar ,
app varchar ,
UNIX_TIMESTAMP(time, 'yyyy-MM-dd HH:mm:ss')*1000 bigint AS xctime ,
WATERMARK FOR xctime AS withOffset( xctime , 1000)
)
WITH (
type='kafka11',
bootstrapServers='kafka:9092',
offsetReset='latest',
topic='test_1',
groupId='flink_sql',
parallelism='4',
timezone='Asia/Shanghai',
topicIsPattern ='false',
sourcedatatype ='dt_nest'
);
CREATE TABLE MyResult(
app VARCHAR,
cnt BIGINT,
wStart timestamp
)WITH(
type ='elasticsearch6',
address ='eshost:9200',
cluster='bigdata-es5.6',
estype ='date',
index ='MyResult',
parallelism ='1',
id='0'
);
insert into MyResult
select
d.app ,
count(d.pushid) cnt ,
TUMBLE_START(d.ROWTIME, INTERVAL '3' SECOND) as wStart
from
MyTable as d
group by d.app,TUMBLE(d.ROWTIME, INTERVAL '3' SECOND);
可否提供一份完整的1.10版本 WATERMARK 语法案例。实际测试中无法使用
- Linguagem predominante
- Java
- Estrelas
- 2k
- Forks
- 913
- Métricas de merge de PRs
- Nenhum PR com merge em 30d
Guia de contribuição
Nenhum guia de contribuição indexado para este repositório
Primeiros passos
- Leia a issue inteira e depois o guia de contribuição do projeto.
- Comente na issue dizendo que vai assumir — evita que duas pessoas façam o mesmo trabalho.
- Faça um fork do repositório e trabalhe em uma branch.
- Abra um pull request que referencie o número da issue.
Mais de DTStack/flinkStreamSQL
-
Dificuldade 5/5 Mais de uma semana Facilidade para iniciantes 20/100
DTStack/flinkStreamSQL#467 · 1 comentário ·
-
Dificuldade 4/5 3-5 dias Facilidade para iniciantes 25/100
DTStack/flinkStreamSQL#437 · 3 comentários ·
-
Dificuldade 3/5 1-2 dias Facilidade para iniciantes 38/100
DTStack/flinkStreamSQL#432 · 1 comentário ·
-
Dificuldade 5/5 Mais de uma semana Facilidade para iniciantes 15/100
DTStack/flinkStreamSQL#431 ·
-
有无对sql 语法校验的方法 Aberta
Dificuldade 5/5 Mais de uma semana Facilidade para iniciantes 15/100
DTStack/flinkStreamSQL#430 · 1 comentário ·
Todas as issues de DTStack/flinkStreamSQL
Issues semelhantes
-
documentation
Dificuldade 2/5 1-3 horas Facilidade para iniciantes 65/100
inu-appcenter/memorIN-backend#288 ·
-
Dificuldade 2/5 1-3 horas Facilidade para iniciantes 65/100
-
frontend maui-pilot pilot-ask question
Dificuldade 2/5 1-3 horas Facilidade para iniciantes 75/100
-
Dificuldade 2/5 1-3 horas Facilidade para iniciantes 75/100
-
executions.Query — startDate and timeRange filters are sent with inverted comparison operators Abertaarea/plugin
Dificuldade 2/5 1-3 horas Facilidade para iniciantes 75/100
kestra-io/plugin-kestra#190 ·