Hacktoberfest 2026: as issues que os mantenedores marcaram para outubro, abertas e boas para iniciantes. Ver issues do Hacktoberfest

1.10版本 无法使用WATERMARK ,报空指针。

Aberta
#334 2 comentários 0 reações 0 responsáveis Ver no GitHub

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
Stack de tecnologia
java, sql

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

  1. Leia a issue inteira e depois o guia de contribuição do projeto.
  2. Comente na issue dizendo que vai assumir — evita que duas pessoas façam o mesmo trabalho.
  3. Faça um fork do repositório e trabalhe em uma branch.
  4. Abra um pull request que referencie o número da issue.

Mais de DTStack/flinkStreamSQL

Todas as issues de DTStack/flinkStreamSQL

Issues semelhantes

Mais issues de Java

Receba novas issues na sua caixa de entrada

Um resumo curto de issues do GitHub para quem está começando.