Commit d55bc00
authored
[Issue 4803][client] return null if the message value/data is not set by producer (#6379)
Fixes #4803
### Motivation
Allow the typed consumer receive messages with `null` value if the producer sends message without payload.
### Modifications
- add a flag in `MessageMetadata` to indicate if the payload is set when the message is created
- check and return `null` if the flag is not set when reading data from a message1 parent 1d5c418 commit d55bc00
File tree
15 files changed
+370
-31
lines changed- pulsar-broker/src
- main/java/org/apache/pulsar/broker/admin/impl
- test/java/org/apache/pulsar/broker
- admin
- service
- pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal
- pulsar-client/src
- main/java/org/apache/pulsar/client/impl
- test/java/org/apache/pulsar/client/impl
- pulsar-common/src/main
- java/org/apache/pulsar/common
- api/proto
- protocol
- proto
- pulsar-flink/src/test/java/org/apache/flink/streaming/connectors/pulsar
- pulsar-functions/worker/src/test/java/org/apache/pulsar/functions/worker
- pulsar-storm/src/test/java/org/apache/pulsar/storm
15 files changed
+370
-31
lines changedLines changed: 3 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
1930 | 1930 | | |
1931 | 1931 | | |
1932 | 1932 | | |
| 1933 | + | |
| 1934 | + | |
| 1935 | + | |
1933 | 1936 | | |
1934 | 1937 | | |
1935 | 1938 | | |
| |||
Lines changed: 34 additions & 14 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
793 | 793 | | |
794 | 794 | | |
795 | 795 | | |
| 796 | + | |
| 797 | + | |
| 798 | + | |
| 799 | + | |
| 800 | + | |
| 801 | + | |
| 802 | + | |
| 803 | + | |
| 804 | + | |
| 805 | + | |
| 806 | + | |
796 | 807 | | |
797 | 808 | | |
798 | 809 | | |
| |||
1559 | 1570 | | |
1560 | 1571 | | |
1561 | 1572 | | |
1562 | | - | |
| 1573 | + | |
1563 | 1574 | | |
1564 | 1575 | | |
1565 | | - | |
| 1576 | + | |
| 1577 | + | |
| 1578 | + | |
| 1579 | + | |
| 1580 | + | |
| 1581 | + | |
1566 | 1582 | | |
1567 | 1583 | | |
1568 | 1584 | | |
1569 | 1585 | | |
1570 | 1586 | | |
1571 | 1587 | | |
1572 | 1588 | | |
1573 | | - | |
1574 | | - | |
| 1589 | + | |
| 1590 | + | |
| 1591 | + | |
| 1592 | + | |
| 1593 | + | |
| 1594 | + | |
1575 | 1595 | | |
1576 | 1596 | | |
1577 | 1597 | | |
| |||
1704 | 1724 | | |
1705 | 1725 | | |
1706 | 1726 | | |
1707 | | - | |
| 1727 | + | |
1708 | 1728 | | |
1709 | 1729 | | |
1710 | 1730 | | |
1711 | 1731 | | |
1712 | 1732 | | |
1713 | | - | |
| 1733 | + | |
1714 | 1734 | | |
1715 | 1735 | | |
1716 | 1736 | | |
| |||
1757 | 1777 | | |
1758 | 1778 | | |
1759 | 1779 | | |
1760 | | - | |
| 1780 | + | |
1761 | 1781 | | |
1762 | 1782 | | |
1763 | 1783 | | |
1764 | 1784 | | |
1765 | | - | |
| 1785 | + | |
1766 | 1786 | | |
1767 | 1787 | | |
1768 | 1788 | | |
1769 | 1789 | | |
1770 | | - | |
| 1790 | + | |
1771 | 1791 | | |
1772 | 1792 | | |
1773 | 1793 | | |
| |||
1829 | 1849 | | |
1830 | 1850 | | |
1831 | 1851 | | |
1832 | | - | |
| 1852 | + | |
1833 | 1853 | | |
1834 | 1854 | | |
1835 | 1855 | | |
1836 | 1856 | | |
1837 | 1857 | | |
1838 | | - | |
| 1858 | + | |
1839 | 1859 | | |
1840 | 1860 | | |
1841 | 1861 | | |
| |||
1866 | 1886 | | |
1867 | 1887 | | |
1868 | 1888 | | |
1869 | | - | |
| 1889 | + | |
1870 | 1890 | | |
1871 | 1891 | | |
1872 | 1892 | | |
| |||
1901 | 1921 | | |
1902 | 1922 | | |
1903 | 1923 | | |
1904 | | - | |
| 1924 | + | |
1905 | 1925 | | |
1906 | 1926 | | |
1907 | 1927 | | |
1908 | | - | |
| 1928 | + | |
1909 | 1929 | | |
1910 | 1930 | | |
1911 | 1931 | | |
| |||
Lines changed: 148 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
| 1 | + | |
| 2 | + | |
| 3 | + | |
| 4 | + | |
| 5 | + | |
| 6 | + | |
| 7 | + | |
| 8 | + | |
| 9 | + | |
| 10 | + | |
| 11 | + | |
| 12 | + | |
| 13 | + | |
| 14 | + | |
| 15 | + | |
| 16 | + | |
| 17 | + | |
| 18 | + | |
| 19 | + | |
| 20 | + | |
| 21 | + | |
| 22 | + | |
| 23 | + | |
| 24 | + | |
| 25 | + | |
| 26 | + | |
| 27 | + | |
| 28 | + | |
| 29 | + | |
| 30 | + | |
| 31 | + | |
| 32 | + | |
| 33 | + | |
| 34 | + | |
| 35 | + | |
| 36 | + | |
| 37 | + | |
| 38 | + | |
| 39 | + | |
| 40 | + | |
| 41 | + | |
| 42 | + | |
| 43 | + | |
| 44 | + | |
| 45 | + | |
| 46 | + | |
| 47 | + | |
| 48 | + | |
| 49 | + | |
| 50 | + | |
| 51 | + | |
| 52 | + | |
| 53 | + | |
| 54 | + | |
| 55 | + | |
| 56 | + | |
| 57 | + | |
| 58 | + | |
| 59 | + | |
| 60 | + | |
| 61 | + | |
| 62 | + | |
| 63 | + | |
| 64 | + | |
| 65 | + | |
| 66 | + | |
| 67 | + | |
| 68 | + | |
| 69 | + | |
| 70 | + | |
| 71 | + | |
| 72 | + | |
| 73 | + | |
| 74 | + | |
| 75 | + | |
| 76 | + | |
| 77 | + | |
| 78 | + | |
| 79 | + | |
| 80 | + | |
| 81 | + | |
| 82 | + | |
| 83 | + | |
| 84 | + | |
| 85 | + | |
| 86 | + | |
| 87 | + | |
| 88 | + | |
| 89 | + | |
| 90 | + | |
| 91 | + | |
| 92 | + | |
| 93 | + | |
| 94 | + | |
| 95 | + | |
| 96 | + | |
| 97 | + | |
| 98 | + | |
| 99 | + | |
| 100 | + | |
| 101 | + | |
| 102 | + | |
| 103 | + | |
| 104 | + | |
| 105 | + | |
| 106 | + | |
| 107 | + | |
| 108 | + | |
| 109 | + | |
| 110 | + | |
| 111 | + | |
| 112 | + | |
| 113 | + | |
| 114 | + | |
| 115 | + | |
| 116 | + | |
| 117 | + | |
| 118 | + | |
| 119 | + | |
| 120 | + | |
| 121 | + | |
| 122 | + | |
| 123 | + | |
| 124 | + | |
| 125 | + | |
| 126 | + | |
| 127 | + | |
| 128 | + | |
| 129 | + | |
| 130 | + | |
| 131 | + | |
| 132 | + | |
| 133 | + | |
| 134 | + | |
| 135 | + | |
| 136 | + | |
| 137 | + | |
| 138 | + | |
| 139 | + | |
| 140 | + | |
| 141 | + | |
| 142 | + | |
| 143 | + | |
| 144 | + | |
| 145 | + | |
| 146 | + | |
| 147 | + | |
| 148 | + | |
Lines changed: 12 additions & 4 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
1288 | 1288 | | |
1289 | 1289 | | |
1290 | 1290 | | |
| 1291 | + | |
1291 | 1292 | | |
1292 | 1293 | | |
1293 | 1294 | | |
| |||
1298 | 1299 | | |
1299 | 1300 | | |
1300 | 1301 | | |
| 1302 | + | |
| 1303 | + | |
| 1304 | + | |
| 1305 | + | |
| 1306 | + | |
| 1307 | + | |
1301 | 1308 | | |
1302 | 1309 | | |
1303 | 1310 | | |
1304 | | - | |
| 1311 | + | |
1305 | 1312 | | |
1306 | 1313 | | |
1307 | 1314 | | |
| |||
1312 | 1319 | | |
1313 | 1320 | | |
1314 | 1321 | | |
1315 | | - | |
| 1322 | + | |
1316 | 1323 | | |
1317 | 1324 | | |
1318 | 1325 | | |
1319 | 1326 | | |
1320 | | - | |
| 1327 | + | |
1321 | 1328 | | |
1322 | 1329 | | |
1323 | 1330 | | |
| |||
1334 | 1341 | | |
1335 | 1342 | | |
1336 | 1343 | | |
1337 | | - | |
| 1344 | + | |
| 1345 | + | |
1338 | 1346 | | |
1339 | 1347 | | |
1340 | 1348 | | |
| |||
Lines changed: 2 additions & 1 deletion
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
557 | 557 | | |
558 | 558 | | |
559 | 559 | | |
560 | | - | |
| 560 | + | |
| 561 | + | |
561 | 562 | | |
562 | 563 | | |
563 | 564 | | |
| |||
Lines changed: 1 addition & 1 deletion
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
1289 | 1289 | | |
1290 | 1290 | | |
1291 | 1291 | | |
1292 | | - | |
| 1292 | + | |
1293 | 1293 | | |
1294 | 1294 | | |
1295 | 1295 | | |
| |||
0 commit comments